Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 3 additions & 0 deletions protocol/messenger.go
Original file line number Diff line number Diff line change
Expand Up @@ -149,6 +149,9 @@ type Messenger struct {
wait chan struct{}
once sync.Once
}
historicSyncMu sync.Mutex
historicSyncInFlight bool
lastHistoricSyncRequestAt time.Time

connectionState connection.State
contractMaker *contracts.ContractMaker
Expand Down
108 changes: 80 additions & 28 deletions protocol/messenger_mailserver.go
Original file line number Diff line number Diff line change
Expand Up @@ -33,6 +33,8 @@ const (
oneMonthDuration = 31 * oneDayDuration

backoffByUserAction = 0 * time.Second

historicSyncMinInterval = 20 * time.Second
)

var ErrNoFiltersForChat = errors.New("no filter registered for given chat")
Expand Down Expand Up @@ -261,8 +263,13 @@ func (m *Messenger) resetFiltersPriority(filters types2.ChatFilters) error {
return nil
}

// RequestAllHistoricMessages requests all the historic messages for any topic
// RequestAllHistoricMessages requests all the historic messages for any topic.
// It keeps aggregating all responses for callers that need the merged payload.
func (m *Messenger) RequestAllHistoricMessages(withRetries bool) (*MessengerResponse, error) {
return m.requestAllHistoricMessages(withRetries, true)
}

func (m *Messenger) requestAllHistoricMessages(withRetries bool, aggregateResponses bool) (*MessengerResponse, error) {
shouldSync, err := m.shouldSync()
if err != nil {
return nil, err
Expand All @@ -276,45 +283,90 @@ func (m *Messenger) RequestAllHistoricMessages(withRetries bool) (*MessengerResp
return nil, nil
}

allResponses := &MessengerResponse{}

filters := m.messaging.ChatFilters()
err = m.updateFiltersPriority(filters)
canSync, err := m.canSyncWithStoreNodes()
if err != nil {
return nil, fmt.Errorf("failed to update filters priority: %w", err)
return nil, err
}
defer func() {
err := m.resetFiltersPriority(filters)
if !canSync {
return nil, nil
}

return m.withHistoricSyncInFlight(time.Now(), func() (*MessengerResponse, error) {
var allResponses *MessengerResponse
if aggregateResponses {
allResponses = &MessengerResponse{}
}

filters := m.messaging.ChatFilters()
err = m.updateFiltersPriority(filters)
if err != nil {
m.logger.Error("failed to reset filters priority", zap.Error(err))
return nil, fmt.Errorf("failed to update filters priority: %w", err)
}
}()
defer func() {
err := m.resetFiltersPriority(filters)
if err != nil {
m.logger.Error("failed to reset filters priority", zap.Error(err))
}
}()

peerInfo := m.messaging.GetActiveStorenode()
peerInfo := m.messaging.GetActiveStorenode()

if withRetries {
response, err := m.performStorenodeTask(func() (*MessengerResponse, error) {
return m.syncFilters(peerInfo, filters)
}, history.WithPeerID(peerInfo.ID))
if withRetries {
response, err := m.performStorenodeTask(func() (*MessengerResponse, error) {
return m.syncFilters(peerInfo, filters)
}, history.WithPeerID(peerInfo.ID))
if err != nil {
return nil, err
}
if aggregateResponses && response != nil {
allResponses.AddChats(response.Chats())
allResponses.AddMessages(response.Messages())
}
return allResponses, nil
}

response, err := m.syncFilters(peerInfo, filters)
if err != nil {
return nil, err
}
if response != nil {
if aggregateResponses && response != nil {
allResponses.AddChats(response.Chats())
allResponses.AddMessages(response.Messages())
}
return allResponses, nil
}
})
}

response, err := m.syncFilters(peerInfo, filters)
if err != nil {
return nil, err
func (m *Messenger) withHistoricSyncInFlight(now time.Time, fn func() (*MessengerResponse, error)) (*MessengerResponse, error) {
m.historicSyncMu.Lock()
if m.historicSyncInFlight {
m.historicSyncMu.Unlock()
m.logger.Debug("skip historic sync request (already in progress)")
return nil, nil
}
if response != nil {
allResponses.AddChats(response.Chats())
allResponses.AddMessages(response.Messages())

if !m.lastHistoricSyncRequestAt.IsZero() {
elapsed := now.Sub(m.lastHistoricSyncRequestAt)
if elapsed < historicSyncMinInterval {
m.historicSyncMu.Unlock()
m.logger.Debug("skip historic sync request (throttled)",
zap.Duration("elapsed", elapsed),
zap.Duration("minInterval", historicSyncMinInterval),
)
return nil, nil
}
}
return allResponses, nil

m.historicSyncInFlight = true
m.lastHistoricSyncRequestAt = now
m.historicSyncMu.Unlock()
defer func() {
m.historicSyncMu.Lock()
m.historicSyncInFlight = false
m.historicSyncMu.Unlock()
}()

return fn()
}

const missingMessageCheckPeriod = 30 * time.Second
Expand Down Expand Up @@ -369,15 +421,15 @@ func (m *Messenger) syncFiltersFrom(peerInfo peer.AddrInfo, filters types2.ChatF
return nil, err
}

topicsData := make(map[string]mailservers.MailserverTopic)
topicsData := make(map[string]mailservers.MailserverTopic, len(topicInfo))
for _, topic := range topicInfo {
topicsData[fmt.Sprintf("%s-%s", topic.PubsubTopic, topic.ContentTopic)] = topic
}

batches := make(map[string]map[int]types2.StoreNodeBatch)

to := m.calculateMailserverTo()
var syncedTopics []mailservers.MailserverTopic
syncedTopics := make([]mailservers.MailserverTopic, 0, len(filters))

sort.Slice(filters[:], func(i, j int) bool {
p1 := filters[i].Priority()
Expand All @@ -396,7 +448,7 @@ func (m *Messenger) syncFiltersFrom(peerInfo peer.AddrInfo, filters types2.ChatF
return nil, err
}

contentTopicsPerPubsubTopic := make(map[string]map[string]*types2.ChatFilter)
contentTopicsPerPubsubTopic := make(map[string]map[string]*types2.ChatFilter, len(filters))
for _, filter := range filters {
if !filter.IsListening() || filter.IsEphemeral() {
continue
Expand Down Expand Up @@ -511,7 +563,7 @@ func (m *Messenger) syncFiltersFrom(peerInfo peer.AddrInfo, filters types2.ChatF
return nil, err
}

var messagesToBeSaved []*common.Message
messagesToBeSaved := make([]*common.Message, 0, len(syncedTopics))
for _, batches := range batches {
for _, batch := range batches {
for _, id := range batch.ChatIDs {
Expand Down
2 changes: 1 addition & 1 deletion protocol/messenger_mailserver_cycle.go
Original file line number Diff line number Diff line change
Expand Up @@ -51,7 +51,7 @@ func (m *Messenger) asyncRequestAllHistoricMessages() {

go func() {
defer gocommon.LogOnPanic()
_, err := m.RequestAllHistoricMessages(true)
_, err := m.requestAllHistoricMessages(true, false)
if err != nil {
m.logger.Error("failed to request historic messages", zap.Error(err))
}
Expand Down
82 changes: 82 additions & 0 deletions protocol/messenger_mailserver_throttle_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,82 @@
package protocol

import (
"errors"
"testing"
"time"

"github.com/stretchr/testify/require"
"go.uber.org/zap"
)

func TestWithHistoricSyncInFlightResetsFlagOnError(t *testing.T) {
m := &Messenger{logger: zap.NewNop()}

_, err := m.withHistoricSyncInFlight(time.Now(), func() (*MessengerResponse, error) {
return nil, errors.New("boom")
})
require.Error(t, err)

m.historicSyncMu.Lock()
require.False(t, m.historicSyncInFlight)
m.historicSyncMu.Unlock()
}

func TestWithHistoricSyncInFlightSkipsWhenAlreadyInFlight(t *testing.T) {
m := &Messenger{logger: zap.NewNop(), historicSyncInFlight: true}

called := false
resp, err := m.withHistoricSyncInFlight(time.Now(), func() (*MessengerResponse, error) {
called = true
return &MessengerResponse{}, nil
})
require.NoError(t, err)
require.Nil(t, resp)
require.False(t, called)

m.historicSyncMu.Lock()
require.True(t, m.historicSyncInFlight)
m.historicSyncMu.Unlock()
}

func TestWithHistoricSyncInFlightSkipsWhenThrottled(t *testing.T) {
now := time.Now()
m := &Messenger{
logger: zap.NewNop(),
lastHistoricSyncRequestAt: now.Add(-(historicSyncMinInterval / 2)),
}

called := false
resp, err := m.withHistoricSyncInFlight(now, func() (*MessengerResponse, error) {
called = true
return &MessengerResponse{}, nil
})
require.NoError(t, err)
require.Nil(t, resp)
require.False(t, called)

m.historicSyncMu.Lock()
require.False(t, m.historicSyncInFlight)
require.Equal(t, now.Add(-(historicSyncMinInterval / 2)), m.lastHistoricSyncRequestAt)
m.historicSyncMu.Unlock()
}

func TestWithHistoricSyncInFlightRunsAndUpdatesTimestamp(t *testing.T) {
now := time.Now()
old := now.Add(-2 * historicSyncMinInterval)
m := &Messenger{logger: zap.NewNop(), lastHistoricSyncRequestAt: old}

called := false
resp, err := m.withHistoricSyncInFlight(now, func() (*MessengerResponse, error) {
called = true
return &MessengerResponse{}, nil
})
require.NoError(t, err)
require.NotNil(t, resp)
require.True(t, called)

m.historicSyncMu.Lock()
require.False(t, m.historicSyncInFlight)
require.Equal(t, now, m.lastHistoricSyncRequestAt)
m.historicSyncMu.Unlock()
}
Loading