Skip to content

Commit 8213ddd

Browse files
committed
test(mailserver): add test that validates throttle behaviour
1 parent 35746a0 commit 8213ddd

2 files changed

Lines changed: 138 additions & 43 deletions

File tree

protocol/messenger_mailserver.go

Lines changed: 56 additions & 43 deletions
Original file line numberDiff line numberDiff line change
@@ -283,7 +283,61 @@ func (m *Messenger) requestAllHistoricMessages(withRetries bool, aggregateRespon
283283
return nil, nil
284284
}
285285

286-
now := time.Now()
286+
canSync, err := m.canSyncWithStoreNodes()
287+
if err != nil {
288+
return nil, err
289+
}
290+
if !canSync {
291+
return nil, nil
292+
}
293+
294+
return m.withHistoricSyncInFlight(time.Now(), func() (*MessengerResponse, error) {
295+
var allResponses *MessengerResponse
296+
if aggregateResponses {
297+
allResponses = &MessengerResponse{}
298+
}
299+
300+
filters := m.messaging.ChatFilters()
301+
err = m.updateFiltersPriority(filters)
302+
if err != nil {
303+
return nil, fmt.Errorf("failed to update filters priority: %w", err)
304+
}
305+
defer func() {
306+
err := m.resetFiltersPriority(filters)
307+
if err != nil {
308+
m.logger.Error("failed to reset filters priority", zap.Error(err))
309+
}
310+
}()
311+
312+
peerInfo := m.messaging.GetActiveStorenode()
313+
314+
if withRetries {
315+
response, err := m.performStorenodeTask(func() (*MessengerResponse, error) {
316+
return m.syncFilters(peerInfo, filters)
317+
}, history.WithPeerID(peerInfo.ID))
318+
if err != nil {
319+
return nil, err
320+
}
321+
if aggregateResponses && response != nil {
322+
allResponses.AddChats(response.Chats())
323+
allResponses.AddMessages(response.Messages())
324+
}
325+
return allResponses, nil
326+
}
327+
328+
response, err := m.syncFilters(peerInfo, filters)
329+
if err != nil {
330+
return nil, err
331+
}
332+
if aggregateResponses && response != nil {
333+
allResponses.AddChats(response.Chats())
334+
allResponses.AddMessages(response.Messages())
335+
}
336+
return allResponses, nil
337+
})
338+
}
339+
340+
func (m *Messenger) withHistoricSyncInFlight(now time.Time, fn func() (*MessengerResponse, error)) (*MessengerResponse, error) {
287341
m.historicSyncMu.Lock()
288342
if m.historicSyncInFlight {
289343
m.historicSyncMu.Unlock()
@@ -312,48 +366,7 @@ func (m *Messenger) requestAllHistoricMessages(withRetries bool, aggregateRespon
312366
m.historicSyncMu.Unlock()
313367
}()
314368

315-
var allResponses *MessengerResponse
316-
if aggregateResponses {
317-
allResponses = &MessengerResponse{}
318-
}
319-
320-
filters := m.messaging.ChatFilters()
321-
err = m.updateFiltersPriority(filters)
322-
if err != nil {
323-
return nil, fmt.Errorf("failed to update filters priority: %w", err)
324-
}
325-
defer func() {
326-
err := m.resetFiltersPriority(filters)
327-
if err != nil {
328-
m.logger.Error("failed to reset filters priority", zap.Error(err))
329-
}
330-
}()
331-
332-
peerInfo := m.messaging.GetActiveStorenode()
333-
334-
if withRetries {
335-
response, err := m.performStorenodeTask(func() (*MessengerResponse, error) {
336-
return m.syncFilters(peerInfo, filters)
337-
}, history.WithPeerID(peerInfo.ID))
338-
if err != nil {
339-
return nil, err
340-
}
341-
if aggregateResponses && response != nil {
342-
allResponses.AddChats(response.Chats())
343-
allResponses.AddMessages(response.Messages())
344-
}
345-
return allResponses, nil
346-
}
347-
348-
response, err := m.syncFilters(peerInfo, filters)
349-
if err != nil {
350-
return nil, err
351-
}
352-
if aggregateResponses && response != nil {
353-
allResponses.AddChats(response.Chats())
354-
allResponses.AddMessages(response.Messages())
355-
}
356-
return allResponses, nil
369+
return fn()
357370
}
358371

359372
const missingMessageCheckPeriod = 30 * time.Second
Lines changed: 82 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,82 @@
1+
package protocol
2+
3+
import (
4+
"errors"
5+
"testing"
6+
"time"
7+
8+
"github.com/stretchr/testify/require"
9+
"go.uber.org/zap"
10+
)
11+
12+
func TestWithHistoricSyncInFlightResetsFlagOnError(t *testing.T) {
13+
m := &Messenger{logger: zap.NewNop()}
14+
15+
_, err := m.withHistoricSyncInFlight(time.Now(), func() (*MessengerResponse, error) {
16+
return nil, errors.New("boom")
17+
})
18+
require.Error(t, err)
19+
20+
m.historicSyncMu.Lock()
21+
require.False(t, m.historicSyncInFlight)
22+
m.historicSyncMu.Unlock()
23+
}
24+
25+
func TestWithHistoricSyncInFlightSkipsWhenAlreadyInFlight(t *testing.T) {
26+
m := &Messenger{logger: zap.NewNop(), historicSyncInFlight: true}
27+
28+
called := false
29+
resp, err := m.withHistoricSyncInFlight(time.Now(), func() (*MessengerResponse, error) {
30+
called = true
31+
return &MessengerResponse{}, nil
32+
})
33+
require.NoError(t, err)
34+
require.Nil(t, resp)
35+
require.False(t, called)
36+
37+
m.historicSyncMu.Lock()
38+
require.True(t, m.historicSyncInFlight)
39+
m.historicSyncMu.Unlock()
40+
}
41+
42+
func TestWithHistoricSyncInFlightSkipsWhenThrottled(t *testing.T) {
43+
now := time.Now()
44+
m := &Messenger{
45+
logger: zap.NewNop(),
46+
lastHistoricSyncRequestAt: now.Add(-(historicSyncMinInterval / 2)),
47+
}
48+
49+
called := false
50+
resp, err := m.withHistoricSyncInFlight(now, func() (*MessengerResponse, error) {
51+
called = true
52+
return &MessengerResponse{}, nil
53+
})
54+
require.NoError(t, err)
55+
require.Nil(t, resp)
56+
require.False(t, called)
57+
58+
m.historicSyncMu.Lock()
59+
require.False(t, m.historicSyncInFlight)
60+
require.Equal(t, now.Add(-(historicSyncMinInterval / 2)), m.lastHistoricSyncRequestAt)
61+
m.historicSyncMu.Unlock()
62+
}
63+
64+
func TestWithHistoricSyncInFlightRunsAndUpdatesTimestamp(t *testing.T) {
65+
now := time.Now()
66+
old := now.Add(-2 * historicSyncMinInterval)
67+
m := &Messenger{logger: zap.NewNop(), lastHistoricSyncRequestAt: old}
68+
69+
called := false
70+
resp, err := m.withHistoricSyncInFlight(now, func() (*MessengerResponse, error) {
71+
called = true
72+
return &MessengerResponse{}, nil
73+
})
74+
require.NoError(t, err)
75+
require.NotNil(t, resp)
76+
require.True(t, called)
77+
78+
m.historicSyncMu.Lock()
79+
require.False(t, m.historicSyncInFlight)
80+
require.Equal(t, now, m.lastHistoricSyncRequestAt)
81+
m.historicSyncMu.Unlock()
82+
}

0 commit comments

Comments
 (0)