-
Notifications
You must be signed in to change notification settings - Fork 951
/
peer_statuses.go
123 lines (109 loc) · 2.71 KB
/
peer_statuses.go
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
// Package peerstatus is a threadsafe global cache to store recent peer status messages for access
// across multiple services.
package peerstatus
import (
"sync"
"time"
"github.com/libp2p/go-libp2p-core/peer"
pb "github.com/prysmaticlabs/prysm/proto/beacon/p2p/v1"
"github.com/prysmaticlabs/prysm/shared/roughtime"
)
var lock sync.RWMutex
var peerStatuses = make(map[peer.ID]*peerStatus)
var failureCount = make(map[peer.ID]int)
var maxFailureThreshold = 3
type peerStatus struct {
status *pb.Status
lastUpdated time.Time
}
// Get most recent status from peer in cache. Threadsafe.
func Get(pid peer.ID) *pb.Status {
lock.RLock()
defer lock.RUnlock()
if pStatus, ok := peerStatuses[pid]; ok {
return pStatus.status
}
return nil
}
// Set most recent status from peer in cache. Threadsafe.
func Set(pid peer.ID, status *pb.Status) {
lock.Lock()
defer lock.Unlock()
if pStatus, ok := peerStatuses[pid]; ok {
pStatus.status = status
peerStatuses[pid] = pStatus
return
}
peerStatuses[pid] = &peerStatus{
status: status,
lastUpdated: roughtime.Now(),
}
}
// Delete peer status from cache. Threadsafe.
func Delete(pid peer.ID) {
lock.Lock()
defer lock.Unlock()
delete(peerStatuses, pid)
}
// Count of peer statuses in cache. Threadsafe.
func Count() int {
lock.RLock()
defer lock.RUnlock()
return len(peerStatuses)
}
// Keys is the list of peer IDs which status exists. Threadsafe.
func Keys() []peer.ID {
lock.RLock()
defer lock.RUnlock()
keys := make([]peer.ID, 0, len(peerStatuses))
for k := range peerStatuses {
keys = append(keys, k)
}
return keys
}
// LastUpdated time which the status was set for the given peer. Threadsafe.
func LastUpdated(pid peer.ID) time.Time {
lock.RLock()
defer lock.RUnlock()
if pStatus, ok := peerStatuses[pid]; ok {
return pStatus.lastUpdated
}
return time.Unix(0, 0)
}
// IncreaseFailureCount increases the failure count for the particular peer.
func IncreaseFailureCount(pid peer.ID) {
lock.Lock()
defer lock.Unlock()
count, ok := failureCount[pid]
if !ok {
failureCount[pid] = 1
return
}
failureCount[pid] = count + 1
}
// FailureCount returns the failure count for the particular peer.
func FailureCount(pid peer.ID) int {
lock.RLock()
defer lock.RUnlock()
count, ok := failureCount[pid]
if !ok {
return 0
}
return count
}
// IsBadPeer checks whether the given peer has
// exceeded the number of bad handshakes threshold.
func IsBadPeer(pid peer.ID) bool {
lock.RLock()
defer lock.RUnlock()
count, ok := failureCount[pid]
if !ok {
return false
}
return count > maxFailureThreshold
}
// Clear the cache. This method should only be used for tests.
func Clear() {
peerStatuses = make(map[peer.ID]*peerStatus)
failureCount = make(map[peer.ID]int)
}