2020-01-13 19:17:30 +00:00
|
|
|
package waku
|
|
|
|
|
|
|
|
import (
|
|
|
|
"bytes"
|
|
|
|
"context"
|
|
|
|
"crypto/ecdsa"
|
|
|
|
"database/sql"
|
|
|
|
"sync"
|
2020-01-20 16:44:32 +00:00
|
|
|
"time"
|
2020-01-13 19:17:30 +00:00
|
|
|
|
|
|
|
"github.com/pkg/errors"
|
|
|
|
"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/transport"
|
|
|
|
)
|
|
|
|
|
|
|
|
var (
|
|
|
|
// ErrNoMailservers returned if there is no configured mailservers that can be used.
|
|
|
|
ErrNoMailservers = errors.New("no configured mailservers")
|
|
|
|
)
|
|
|
|
|
|
|
|
type wakuServiceKeysManager struct {
|
|
|
|
waku types.Waku
|
|
|
|
|
|
|
|
// Identity of the current user.
|
|
|
|
privateKey *ecdsa.PrivateKey
|
|
|
|
|
|
|
|
passToSymKeyMutex sync.RWMutex
|
|
|
|
passToSymKeyCache map[string]string
|
|
|
|
}
|
|
|
|
|
|
|
|
func (m *wakuServiceKeysManager) AddOrGetKeyPair(priv *ecdsa.PrivateKey) (string, error) {
|
|
|
|
// caching is handled in waku
|
|
|
|
return m.waku.AddKeyPair(priv)
|
|
|
|
}
|
|
|
|
|
|
|
|
func (m *wakuServiceKeysManager) AddOrGetSymKeyFromPassword(password string) (string, error) {
|
|
|
|
m.passToSymKeyMutex.Lock()
|
|
|
|
defer m.passToSymKeyMutex.Unlock()
|
|
|
|
|
|
|
|
if val, ok := m.passToSymKeyCache[password]; ok {
|
|
|
|
return val, nil
|
|
|
|
}
|
|
|
|
|
|
|
|
id, err := m.waku.AddSymKeyFromPassword(password)
|
|
|
|
if err != nil {
|
|
|
|
return id, err
|
|
|
|
}
|
|
|
|
|
|
|
|
m.passToSymKeyCache[password] = id
|
|
|
|
|
|
|
|
return id, nil
|
|
|
|
}
|
|
|
|
|
|
|
|
func (m *wakuServiceKeysManager) RawSymKey(id string) ([]byte, error) {
|
|
|
|
return m.waku.GetSymKey(id)
|
|
|
|
}
|
|
|
|
|
2020-02-10 11:22:37 +00:00
|
|
|
type Option func(*Transport) error
|
2020-01-13 19:17:30 +00:00
|
|
|
|
2020-02-10 11:22:37 +00:00
|
|
|
// Transport is a transport based on Whisper service.
|
|
|
|
type Transport struct {
|
2020-01-13 19:17:30 +00:00
|
|
|
waku types.Waku
|
|
|
|
api types.PublicWakuAPI // only PublicWakuAPI implements logic to send messages
|
|
|
|
keysManager *wakuServiceKeysManager
|
|
|
|
filters *transport.FiltersManager
|
|
|
|
logger *zap.Logger
|
2021-01-08 15:21:25 +00:00
|
|
|
cache *transport.ProcessedMessageIDsCache
|
2020-01-13 19:17:30 +00:00
|
|
|
|
|
|
|
mailservers []string
|
|
|
|
envelopesMonitor *EnvelopesMonitor
|
2021-01-14 22:15:13 +00:00
|
|
|
quit chan struct{}
|
2020-01-13 19:17:30 +00:00
|
|
|
}
|
|
|
|
|
2020-02-10 11:22:37 +00:00
|
|
|
// NewTransport returns a new Transport.
|
2020-01-13 19:17:30 +00:00
|
|
|
// TODO: leaving a chat should verify that for a given public key
|
|
|
|
// there are no other chats. It may happen that we leave a private chat
|
|
|
|
// but still have a public chat for a given public key.
|
2020-02-10 11:22:37 +00:00
|
|
|
func NewTransport(
|
2020-01-13 19:17:30 +00:00
|
|
|
waku types.Waku,
|
|
|
|
privateKey *ecdsa.PrivateKey,
|
|
|
|
db *sql.DB,
|
|
|
|
mailservers []string,
|
|
|
|
envelopesMonitorConfig *transport.EnvelopesMonitorConfig,
|
|
|
|
logger *zap.Logger,
|
|
|
|
opts ...Option,
|
2020-02-10 11:22:37 +00:00
|
|
|
) (*Transport, error) {
|
2020-01-13 19:17:30 +00:00
|
|
|
filtersManager, err := transport.NewFiltersManager(newSQLitePersistence(db), waku, privateKey, logger)
|
|
|
|
if err != nil {
|
|
|
|
return nil, err
|
|
|
|
}
|
|
|
|
|
|
|
|
var envelopesMonitor *EnvelopesMonitor
|
|
|
|
if envelopesMonitorConfig != nil {
|
|
|
|
envelopesMonitor = NewEnvelopesMonitor(waku, *envelopesMonitorConfig)
|
|
|
|
envelopesMonitor.Start()
|
|
|
|
}
|
|
|
|
|
|
|
|
var api types.PublicWhisperAPI
|
|
|
|
if waku != nil {
|
|
|
|
api = waku.PublicWakuAPI()
|
|
|
|
}
|
2020-02-10 11:22:37 +00:00
|
|
|
t := &Transport{
|
2020-01-13 19:17:30 +00:00
|
|
|
waku: waku,
|
|
|
|
api: api,
|
2021-01-08 15:21:25 +00:00
|
|
|
cache: transport.NewProcessedMessageIDsCache(db),
|
2020-01-13 19:17:30 +00:00
|
|
|
envelopesMonitor: envelopesMonitor,
|
2021-01-14 22:15:13 +00:00
|
|
|
quit: make(chan struct{}),
|
2020-01-13 19:17:30 +00:00
|
|
|
keysManager: &wakuServiceKeysManager{
|
|
|
|
waku: waku,
|
|
|
|
privateKey: privateKey,
|
|
|
|
passToSymKeyCache: make(map[string]string),
|
|
|
|
},
|
|
|
|
filters: filtersManager,
|
|
|
|
mailservers: mailservers,
|
2020-02-10 11:22:37 +00:00
|
|
|
logger: logger.With(zap.Namespace("Transport")),
|
2020-01-13 19:17:30 +00:00
|
|
|
}
|
|
|
|
|
|
|
|
for _, opt := range opts {
|
|
|
|
if err := opt(t); err != nil {
|
|
|
|
return nil, err
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
2021-01-14 22:15:13 +00:00
|
|
|
t.cleanFiltersLoop()
|
|
|
|
|
2020-01-13 19:17:30 +00:00
|
|
|
return t, nil
|
|
|
|
}
|
|
|
|
|
2020-02-10 11:22:37 +00:00
|
|
|
func (a *Transport) InitFilters(chatIDs []string, publicKeys []*ecdsa.PublicKey) ([]*transport.Filter, error) {
|
2020-01-13 19:17:30 +00:00
|
|
|
return a.filters.Init(chatIDs, publicKeys)
|
|
|
|
}
|
|
|
|
|
2020-11-18 09:16:51 +00:00
|
|
|
func (a *Transport) InitPublicFilters(chatIDs []string) ([]*transport.Filter, error) {
|
|
|
|
return a.filters.InitPublicFilters(chatIDs)
|
|
|
|
}
|
|
|
|
|
2020-02-10 11:22:37 +00:00
|
|
|
func (a *Transport) Filters() []*transport.Filter {
|
2020-01-13 19:17:30 +00:00
|
|
|
return a.filters.Filters()
|
|
|
|
}
|
|
|
|
|
2020-02-10 11:22:37 +00:00
|
|
|
func (a *Transport) LoadFilters(filters []*transport.Filter) ([]*transport.Filter, error) {
|
2020-01-13 19:17:30 +00:00
|
|
|
return a.filters.InitWithFilters(filters)
|
|
|
|
}
|
|
|
|
|
2021-01-11 10:32:51 +00:00
|
|
|
func (a *Transport) InitCommunityFilters(pks []*ecdsa.PrivateKey) ([]*transport.Filter, error) {
|
|
|
|
return a.filters.InitCommunityFilters(pks)
|
|
|
|
}
|
|
|
|
|
2020-02-10 11:22:37 +00:00
|
|
|
func (a *Transport) RemoveFilters(filters []*transport.Filter) error {
|
2020-01-13 19:17:30 +00:00
|
|
|
return a.filters.Remove(filters...)
|
|
|
|
}
|
|
|
|
|
2020-12-22 17:12:03 +00:00
|
|
|
func (a *Transport) RemoveFilterByChatID(chatID string) (*transport.Filter, error) {
|
2020-11-18 09:16:51 +00:00
|
|
|
return a.filters.RemoveFilterByChatID(chatID)
|
|
|
|
}
|
|
|
|
|
2020-02-10 11:22:37 +00:00
|
|
|
func (a *Transport) ResetFilters() error {
|
2020-01-13 19:17:30 +00:00
|
|
|
return a.filters.Reset()
|
|
|
|
}
|
|
|
|
|
2020-02-10 11:22:37 +00:00
|
|
|
func (a *Transport) ProcessNegotiatedSecret(secret types.NegotiatedSecret) (*transport.Filter, error) {
|
2020-01-13 19:17:30 +00:00
|
|
|
filter, err := a.filters.LoadNegotiated(secret)
|
|
|
|
if err != nil {
|
|
|
|
return nil, err
|
|
|
|
}
|
|
|
|
return filter, nil
|
|
|
|
}
|
|
|
|
|
2021-01-11 10:32:51 +00:00
|
|
|
func (a *Transport) JoinPublic(chatID string) (*transport.Filter, error) {
|
|
|
|
return a.filters.LoadPublic(chatID)
|
2020-01-13 19:17:30 +00:00
|
|
|
}
|
|
|
|
|
2020-02-10 11:22:37 +00:00
|
|
|
func (a *Transport) LeavePublic(chatID string) error {
|
2020-01-13 19:17:30 +00:00
|
|
|
chat := a.filters.Filter(chatID)
|
|
|
|
if chat != nil {
|
|
|
|
return nil
|
|
|
|
}
|
|
|
|
return a.filters.Remove(chat)
|
|
|
|
}
|
|
|
|
|
2021-01-11 10:32:51 +00:00
|
|
|
func (a *Transport) JoinPrivate(publicKey *ecdsa.PublicKey) (*transport.Filter, error) {
|
|
|
|
return a.filters.LoadContactCode(publicKey)
|
2020-01-13 19:17:30 +00:00
|
|
|
}
|
|
|
|
|
2020-02-10 11:22:37 +00:00
|
|
|
func (a *Transport) LeavePrivate(publicKey *ecdsa.PublicKey) error {
|
2020-01-13 19:17:30 +00:00
|
|
|
filters := a.filters.FiltersByPublicKey(publicKey)
|
|
|
|
return a.filters.Remove(filters...)
|
|
|
|
}
|
|
|
|
|
2021-01-11 10:32:51 +00:00
|
|
|
func (a *Transport) JoinGroup(publicKeys []*ecdsa.PublicKey) ([]*transport.Filter, error) {
|
|
|
|
var filters []*transport.Filter
|
2020-01-13 19:17:30 +00:00
|
|
|
for _, pk := range publicKeys {
|
2021-01-11 10:32:51 +00:00
|
|
|
f, err := a.filters.LoadContactCode(pk)
|
2020-01-13 19:17:30 +00:00
|
|
|
if err != nil {
|
2021-01-11 10:32:51 +00:00
|
|
|
return nil, err
|
2020-01-13 19:17:30 +00:00
|
|
|
}
|
2021-01-11 10:32:51 +00:00
|
|
|
filters = append(filters, f)
|
|
|
|
|
2020-01-13 19:17:30 +00:00
|
|
|
}
|
2021-01-11 10:32:51 +00:00
|
|
|
return filters, nil
|
2020-01-13 19:17:30 +00:00
|
|
|
}
|
|
|
|
|
2020-02-10 11:22:37 +00:00
|
|
|
func (a *Transport) LeaveGroup(publicKeys []*ecdsa.PublicKey) error {
|
2020-01-13 19:17:30 +00:00
|
|
|
for _, publicKey := range publicKeys {
|
|
|
|
filters := a.filters.FiltersByPublicKey(publicKey)
|
|
|
|
if err := a.filters.Remove(filters...); err != nil {
|
|
|
|
return err
|
|
|
|
}
|
|
|
|
}
|
|
|
|
return nil
|
|
|
|
}
|
|
|
|
|
2020-02-10 11:22:37 +00:00
|
|
|
func (a *Transport) RetrieveRawAll() (map[transport.Filter][]*types.Message, error) {
|
2020-01-13 19:17:30 +00:00
|
|
|
result := make(map[transport.Filter][]*types.Message)
|
|
|
|
|
|
|
|
allFilters := a.filters.Filters()
|
|
|
|
for _, filter := range allFilters {
|
2021-01-14 22:15:13 +00:00
|
|
|
// Don't pull from filters we don't listen to
|
|
|
|
if !filter.Listen {
|
|
|
|
continue
|
|
|
|
}
|
2020-01-13 19:17:30 +00:00
|
|
|
msgs, err := a.api.GetFilterMessages(filter.FilterID)
|
|
|
|
if err != nil {
|
2021-01-08 15:21:25 +00:00
|
|
|
a.logger.Warn("failed to fetch messages", zap.Error(err))
|
2020-01-13 19:17:30 +00:00
|
|
|
continue
|
|
|
|
}
|
2021-01-08 15:21:25 +00:00
|
|
|
if len(msgs) == 0 {
|
|
|
|
continue
|
|
|
|
}
|
|
|
|
|
|
|
|
ids := make([]string, len(msgs))
|
|
|
|
for i := range msgs {
|
|
|
|
id := types.EncodeHex(msgs[i].Hash)
|
|
|
|
ids[i] = id
|
|
|
|
}
|
|
|
|
|
|
|
|
hits, err := a.cache.Hits(ids)
|
|
|
|
if err != nil {
|
|
|
|
a.logger.Error("failed to check messages exists", zap.Error(err))
|
|
|
|
return nil, err
|
|
|
|
}
|
|
|
|
|
|
|
|
for i := range msgs {
|
|
|
|
// Exclude anything that is a cache hit
|
|
|
|
if !hits[types.EncodeHex(msgs[i].Hash)] {
|
|
|
|
result[*filter] = append(result[*filter], msgs[i])
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
2020-01-13 19:17:30 +00:00
|
|
|
}
|
|
|
|
|
|
|
|
return result, nil
|
|
|
|
}
|
|
|
|
|
|
|
|
// SendPublic sends a new message using the Whisper service.
|
|
|
|
// For public filters, chat name is used as an ID as well as
|
|
|
|
// a topic.
|
2020-02-10 11:22:37 +00:00
|
|
|
func (a *Transport) SendPublic(ctx context.Context, newMessage *types.NewMessage, chatName string) ([]byte, error) {
|
2020-01-13 19:17:30 +00:00
|
|
|
if err := a.addSig(newMessage); err != nil {
|
|
|
|
return nil, err
|
|
|
|
}
|
|
|
|
|
|
|
|
filter, err := a.filters.LoadPublic(chatName)
|
|
|
|
if err != nil {
|
|
|
|
return nil, err
|
|
|
|
}
|
|
|
|
|
|
|
|
newMessage.SymKeyID = filter.SymKeyID
|
|
|
|
newMessage.Topic = filter.Topic
|
|
|
|
|
|
|
|
return a.api.Post(ctx, *newMessage)
|
|
|
|
}
|
|
|
|
|
2020-02-10 11:22:37 +00:00
|
|
|
func (a *Transport) SendPrivateWithSharedSecret(ctx context.Context, newMessage *types.NewMessage, publicKey *ecdsa.PublicKey, secret []byte) ([]byte, error) {
|
2020-01-13 19:17:30 +00:00
|
|
|
if err := a.addSig(newMessage); err != nil {
|
|
|
|
return nil, err
|
|
|
|
}
|
|
|
|
|
|
|
|
filter, err := a.filters.LoadNegotiated(types.NegotiatedSecret{
|
|
|
|
PublicKey: publicKey,
|
|
|
|
Key: secret,
|
|
|
|
})
|
|
|
|
if err != nil {
|
|
|
|
return nil, err
|
|
|
|
}
|
|
|
|
|
|
|
|
newMessage.SymKeyID = filter.SymKeyID
|
|
|
|
newMessage.Topic = filter.Topic
|
|
|
|
newMessage.PublicKey = nil
|
|
|
|
|
|
|
|
return a.api.Post(ctx, *newMessage)
|
|
|
|
}
|
|
|
|
|
2020-02-10 11:22:37 +00:00
|
|
|
func (a *Transport) SendPrivateWithPartitioned(ctx context.Context, newMessage *types.NewMessage, publicKey *ecdsa.PublicKey) ([]byte, error) {
|
2020-01-13 19:17:30 +00:00
|
|
|
if err := a.addSig(newMessage); err != nil {
|
|
|
|
return nil, err
|
|
|
|
}
|
|
|
|
|
2020-07-21 15:41:10 +00:00
|
|
|
filter, err := a.filters.LoadPartitioned(publicKey, a.keysManager.privateKey, false)
|
2020-01-13 19:17:30 +00:00
|
|
|
if err != nil {
|
|
|
|
return nil, err
|
|
|
|
}
|
|
|
|
|
|
|
|
newMessage.Topic = filter.Topic
|
|
|
|
newMessage.PublicKey = crypto.FromECDSAPub(publicKey)
|
|
|
|
|
|
|
|
return a.api.Post(ctx, *newMessage)
|
|
|
|
}
|
|
|
|
|
2021-01-18 09:12:03 +00:00
|
|
|
func (a *Transport) SendPrivateOnPersonalTopic(ctx context.Context, newMessage *types.NewMessage, publicKey *ecdsa.PublicKey) ([]byte, error) {
|
|
|
|
if err := a.addSig(newMessage); err != nil {
|
|
|
|
return nil, err
|
|
|
|
}
|
|
|
|
|
|
|
|
filter, err := a.filters.LoadPersonal(publicKey, a.keysManager.privateKey, false)
|
|
|
|
if err != nil {
|
|
|
|
return nil, err
|
|
|
|
}
|
|
|
|
|
|
|
|
newMessage.Topic = filter.Topic
|
|
|
|
newMessage.PublicKey = crypto.FromECDSAPub(publicKey)
|
|
|
|
|
|
|
|
return a.api.Post(ctx, *newMessage)
|
|
|
|
}
|
|
|
|
|
2020-07-21 15:41:10 +00:00
|
|
|
func (a *Transport) LoadKeyFilters(key *ecdsa.PrivateKey) (*transport.Filter, error) {
|
2021-01-14 22:15:13 +00:00
|
|
|
return a.filters.LoadEphemeral(&key.PublicKey, key, true)
|
2020-07-21 15:41:10 +00:00
|
|
|
}
|
|
|
|
|
2021-01-11 10:32:51 +00:00
|
|
|
func (a *Transport) SendCommunityMessage(ctx context.Context, newMessage *types.NewMessage, publicKey *ecdsa.PublicKey) ([]byte, error) {
|
2020-01-13 19:17:30 +00:00
|
|
|
if err := a.addSig(newMessage); err != nil {
|
|
|
|
return nil, err
|
|
|
|
}
|
|
|
|
|
2021-01-11 10:32:51 +00:00
|
|
|
// We load the filter to make sure we can post on it
|
|
|
|
filter, err := a.filters.LoadPublic(transport.PubkeyToHex(publicKey)[2:])
|
|
|
|
if err != nil {
|
|
|
|
return nil, err
|
|
|
|
}
|
2020-01-13 19:17:30 +00:00
|
|
|
|
2021-01-11 10:32:51 +00:00
|
|
|
newMessage.Topic = filter.Topic
|
2020-01-13 19:17:30 +00:00
|
|
|
newMessage.PublicKey = crypto.FromECDSAPub(publicKey)
|
|
|
|
|
2021-01-11 10:32:51 +00:00
|
|
|
a.logger.Debug("SENDING message", zap.Binary("topic", filter.Topic[:]))
|
|
|
|
|
2020-01-13 19:17:30 +00:00
|
|
|
return a.api.Post(ctx, *newMessage)
|
|
|
|
}
|
|
|
|
|
2021-01-14 22:15:13 +00:00
|
|
|
func (a *Transport) cleanFilters() error {
|
|
|
|
return a.filters.RemoveNoListenFilters()
|
|
|
|
}
|
|
|
|
|
2020-02-10 11:22:37 +00:00
|
|
|
func (a *Transport) addSig(newMessage *types.NewMessage) error {
|
2020-01-13 19:17:30 +00:00
|
|
|
sigID, err := a.keysManager.AddOrGetKeyPair(a.keysManager.privateKey)
|
|
|
|
if err != nil {
|
|
|
|
return err
|
|
|
|
}
|
|
|
|
newMessage.SigID = sigID
|
|
|
|
return nil
|
|
|
|
}
|
|
|
|
|
2020-02-10 11:22:37 +00:00
|
|
|
func (a *Transport) Track(identifiers [][]byte, hash []byte, newMessage *types.NewMessage) {
|
2020-01-13 19:17:30 +00:00
|
|
|
if a.envelopesMonitor != nil {
|
|
|
|
a.envelopesMonitor.Add(identifiers, types.BytesToHash(hash), *newMessage)
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
2020-01-20 16:44:32 +00:00
|
|
|
// GetCurrentTime returns the current unix timestamp in milliseconds
|
2020-02-10 11:22:37 +00:00
|
|
|
func (a *Transport) GetCurrentTime() uint64 {
|
2020-01-20 16:44:32 +00:00
|
|
|
return uint64(a.waku.GetCurrentTime().UnixNano() / int64(time.Millisecond))
|
|
|
|
}
|
|
|
|
|
2020-11-03 12:42:42 +00:00
|
|
|
func (a *Transport) MaxMessageSize() uint32 {
|
|
|
|
return a.waku.MaxMessageSize()
|
|
|
|
}
|
|
|
|
|
2020-02-10 11:22:37 +00:00
|
|
|
func (a *Transport) Stop() error {
|
2021-01-14 22:15:13 +00:00
|
|
|
close(a.quit)
|
2020-01-13 19:17:30 +00:00
|
|
|
if a.envelopesMonitor != nil {
|
|
|
|
a.envelopesMonitor.Stop()
|
|
|
|
}
|
|
|
|
return nil
|
|
|
|
}
|
|
|
|
|
2021-01-14 22:15:13 +00:00
|
|
|
// cleanFiltersLoop cleans up the topic we create for the only purpose
|
|
|
|
// of sending messages.
|
|
|
|
// Whenever we send a message we also need to listen to that particular topic
|
|
|
|
// but in case of asymettric topics, we are not interested in listening to them.
|
|
|
|
// We therefore periodically clean them up so we don't receive unnecessary data.
|
|
|
|
|
|
|
|
func (a *Transport) cleanFiltersLoop() {
|
|
|
|
|
|
|
|
ticker := time.NewTicker(5 * time.Minute)
|
|
|
|
go func() {
|
|
|
|
for {
|
|
|
|
select {
|
|
|
|
case <-a.quit:
|
|
|
|
ticker.Stop()
|
|
|
|
return
|
|
|
|
case <-ticker.C:
|
|
|
|
err := a.cleanFilters()
|
|
|
|
if err != nil {
|
|
|
|
a.logger.Error("failed to clean up topics", zap.Error(err))
|
|
|
|
}
|
|
|
|
}
|
|
|
|
}
|
|
|
|
}()
|
|
|
|
}
|
|
|
|
|
2020-01-13 19:17:30 +00:00
|
|
|
// RequestHistoricMessages requests historic messages for all registered filters.
|
2020-02-10 11:22:37 +00:00
|
|
|
func (a *Transport) SendMessagesRequest(
|
2020-01-13 19:17:30 +00:00
|
|
|
ctx context.Context,
|
|
|
|
peerID []byte,
|
|
|
|
from, to uint32,
|
|
|
|
previousCursor []byte,
|
|
|
|
) (cursor []byte, err error) {
|
|
|
|
topics := make([]types.TopicType, len(a.Filters()))
|
|
|
|
for _, f := range a.Filters() {
|
|
|
|
topics = append(topics, f.Topic)
|
|
|
|
}
|
|
|
|
|
|
|
|
r := createMessagesRequest(from, to, previousCursor, topics)
|
|
|
|
r.SetDefaults(a.waku.GetCurrentTime())
|
|
|
|
|
|
|
|
events := make(chan types.EnvelopeEvent, 10)
|
|
|
|
sub := a.waku.SubscribeEnvelopeEvents(events)
|
|
|
|
defer sub.Unsubscribe()
|
|
|
|
|
|
|
|
err = a.waku.SendMessagesRequest(peerID, r)
|
|
|
|
if err != nil {
|
|
|
|
return
|
|
|
|
}
|
|
|
|
|
|
|
|
resp, err := a.waitForRequestCompleted(ctx, r.ID, events)
|
|
|
|
if err == nil && resp != nil && resp.Error != nil {
|
|
|
|
err = resp.Error
|
|
|
|
} else if err == nil && resp != nil {
|
|
|
|
cursor = resp.Cursor
|
|
|
|
}
|
|
|
|
return
|
|
|
|
}
|
|
|
|
|
2020-02-10 11:22:37 +00:00
|
|
|
func (a *Transport) waitForRequestCompleted(ctx context.Context, requestID []byte, events chan types.EnvelopeEvent) (*types.MailServerResponse, error) {
|
2020-01-13 19:17:30 +00:00
|
|
|
for {
|
|
|
|
select {
|
|
|
|
case ev := <-events:
|
|
|
|
a.logger.Debug(
|
|
|
|
"waiting for request completed and received an event",
|
|
|
|
zap.Binary("requestID", requestID),
|
|
|
|
zap.Any("event", ev),
|
|
|
|
)
|
|
|
|
if !bytes.Equal(ev.Hash.Bytes(), requestID) {
|
|
|
|
continue
|
|
|
|
}
|
|
|
|
if ev.Event != types.EventMailServerRequestCompleted {
|
|
|
|
continue
|
|
|
|
}
|
|
|
|
data, ok := ev.Data.(*types.MailServerResponse)
|
|
|
|
if ok {
|
|
|
|
return data, nil
|
|
|
|
}
|
|
|
|
case <-ctx.Done():
|
|
|
|
return nil, ctx.Err()
|
|
|
|
}
|
|
|
|
}
|
|
|
|
}
|
2020-12-15 14:43:41 +00:00
|
|
|
|
2021-01-08 15:21:25 +00:00
|
|
|
// ConfirmMessagesProcessed marks the messages as processed in the cache so
|
|
|
|
// they won't be passed to the next layer anymore
|
|
|
|
func (a *Transport) ConfirmMessagesProcessed(ids []string, timestamp uint64) error {
|
|
|
|
return a.cache.Add(ids, timestamp)
|
|
|
|
}
|
|
|
|
|
|
|
|
// CleanMessagesProcessed clears the messages that are older than timestamp
|
|
|
|
func (a *Transport) CleanMessagesProcessed(timestamp uint64) error {
|
|
|
|
return a.cache.Clean(timestamp)
|
|
|
|
}
|
|
|
|
|
2020-12-15 14:43:41 +00:00
|
|
|
func (a *Transport) SetEnvelopeEventsHandler(handler transport.EnvelopeEventsHandler) error {
|
|
|
|
if a.envelopesMonitor == nil {
|
|
|
|
return errors.New("Current transport has no envelopes monitor")
|
|
|
|
}
|
|
|
|
a.envelopesMonitor.handler = handler
|
|
|
|
return nil
|
|
|
|
}
|