master
go 240 lines 6.11 KB
Raw
1 package libp2p
2
3 import (
4 "context"
5 "errors"
6 "net"
7 "sync"
8 "time"
9
10 "github.com/libp2p/go-libp2p/core/network"
11 "github.com/libp2p/go-libp2p/core/peer"
12 "github.com/libp2p/go-libp2p/core/protocol"
13 rcmgr "github.com/libp2p/go-libp2p/p2p/host/resource-manager"
14 ma "github.com/multiformats/go-multiaddr"
15 "go.uber.org/zap"
16 )
17
18 type loggingResourceManager struct {
19 logger *zap.SugaredLogger
20 delegate network.ResourceManager
21 logInterval time.Duration
22
23 mut sync.Mutex
24 limitExceededErrs map[string]int
25 }
26
27 type loggingScope struct {
28 logger *zap.SugaredLogger
29 delegate network.ResourceScope
30 countErrs func(error)
31 }
32
33 var (
34 _ network.ResourceManager = (*loggingResourceManager)(nil)
35 _ rcmgr.ResourceManagerState = (*loggingResourceManager)(nil)
36 )
37
38 func (n *loggingResourceManager) start(ctx context.Context) {
39 logInterval := n.logInterval
40 if logInterval == 0 {
41 logInterval = 10 * time.Second
42 }
43 ticker := time.NewTicker(logInterval)
44 go func() {
45 defer ticker.Stop()
46 for {
47 select {
48 case <-ticker.C:
49 n.mut.Lock()
50 errs := n.limitExceededErrs
51 n.limitExceededErrs = make(map[string]int)
52
53 for e, count := range errs {
54 n.logger.Warnf("Protected from exceeding resource limits %d times. libp2p message: %q.", count, e)
55 }
56
57 if len(errs) != 0 {
58 n.logger.Warnf("Learn more about potential actions to take at: https://github.com/ipfs/kubo/blob/master/docs/libp2p-resource-management.md")
59 }
60
61 n.mut.Unlock()
62 case <-ctx.Done():
63 return
64 }
65 }
66 }()
67 }
68
69 func (n *loggingResourceManager) countErrs(err error) {
70 if errors.Is(err, network.ErrResourceLimitExceeded) {
71 n.mut.Lock()
72 if n.limitExceededErrs == nil {
73 n.limitExceededErrs = make(map[string]int)
74 }
75
76 // we need to unwrap the error to get the limit scope and the kind of reached limit
77 eout := errors.Unwrap(err)
78 if eout != nil {
79 n.limitExceededErrs[eout.Error()]++
80 }
81
82 n.mut.Unlock()
83 }
84 }
85
86 func (n *loggingResourceManager) ViewSystem(f func(network.ResourceScope) error) error {
87 return n.delegate.ViewSystem(f)
88 }
89
90 func (n *loggingResourceManager) ViewTransient(f func(network.ResourceScope) error) error {
91 return n.delegate.ViewTransient(func(s network.ResourceScope) error {
92 return f(&loggingScope{logger: n.logger, delegate: s, countErrs: n.countErrs})
93 })
94 }
95
96 func (n *loggingResourceManager) ViewService(svc string, f func(network.ServiceScope) error) error {
97 return n.delegate.ViewService(svc, func(s network.ServiceScope) error {
98 return f(&loggingScope{logger: n.logger, delegate: s, countErrs: n.countErrs})
99 })
100 }
101
102 func (n *loggingResourceManager) ViewProtocol(p protocol.ID, f func(network.ProtocolScope) error) error {
103 return n.delegate.ViewProtocol(p, func(s network.ProtocolScope) error {
104 return f(&loggingScope{logger: n.logger, delegate: s, countErrs: n.countErrs})
105 })
106 }
107
108 func (n *loggingResourceManager) ViewPeer(p peer.ID, f func(network.PeerScope) error) error {
109 return n.delegate.ViewPeer(p, func(s network.PeerScope) error {
110 return f(&loggingScope{logger: n.logger, delegate: s, countErrs: n.countErrs})
111 })
112 }
113
114 func (n *loggingResourceManager) OpenConnection(dir network.Direction, usefd bool, remote ma.Multiaddr) (network.ConnManagementScope, error) {
115 connMgmtScope, err := n.delegate.OpenConnection(dir, usefd, remote)
116 n.countErrs(err)
117 return connMgmtScope, err
118 }
119
120 func (n *loggingResourceManager) OpenStream(p peer.ID, dir network.Direction) (network.StreamManagementScope, error) {
121 connMgmtScope, err := n.delegate.OpenStream(p, dir)
122 n.countErrs(err)
123 return connMgmtScope, err
124 }
125
126 func (n *loggingResourceManager) Close() error {
127 return n.delegate.Close()
128 }
129
130 func (n *loggingResourceManager) ListServices() []string {
131 rapi, ok := n.delegate.(rcmgr.ResourceManagerState)
132 if !ok {
133 return nil
134 }
135
136 return rapi.ListServices()
137 }
138
139 func (n *loggingResourceManager) ListProtocols() []protocol.ID {
140 rapi, ok := n.delegate.(rcmgr.ResourceManagerState)
141 if !ok {
142 return nil
143 }
144
145 return rapi.ListProtocols()
146 }
147
148 func (n *loggingResourceManager) ListPeers() []peer.ID {
149 rapi, ok := n.delegate.(rcmgr.ResourceManagerState)
150 if !ok {
151 return nil
152 }
153
154 return rapi.ListPeers()
155 }
156
157 func (n *loggingResourceManager) Stat() rcmgr.ResourceManagerStat {
158 rapi, ok := n.delegate.(rcmgr.ResourceManagerState)
159 if !ok {
160 return rcmgr.ResourceManagerStat{}
161 }
162
163 return rapi.Stat()
164 }
165
166 func (n *loggingResourceManager) VerifySourceAddress(addr net.Addr) bool {
167 return n.delegate.VerifySourceAddress(addr)
168 }
169
170 func (s *loggingScope) ReserveMemory(size int, prio uint8) error {
171 err := s.delegate.ReserveMemory(size, prio)
172 s.countErrs(err)
173 return err
174 }
175
176 func (s *loggingScope) ReleaseMemory(size int) {
177 s.delegate.ReleaseMemory(size)
178 }
179
180 func (s *loggingScope) Stat() network.ScopeStat {
181 return s.delegate.Stat()
182 }
183
184 func (s *loggingScope) BeginSpan() (network.ResourceScopeSpan, error) {
185 return s.delegate.BeginSpan()
186 }
187
188 func (s *loggingScope) Done() {
189 s.delegate.(network.ResourceScopeSpan).Done()
190 }
191
192 func (s *loggingScope) Name() string {
193 return s.delegate.(network.ServiceScope).Name()
194 }
195
196 func (s *loggingScope) Protocol() protocol.ID {
197 return s.delegate.(network.ProtocolScope).Protocol()
198 }
199
200 func (s *loggingScope) Peer() peer.ID {
201 return s.delegate.(network.PeerScope).Peer()
202 }
203
204 func (s *loggingScope) PeerScope() network.PeerScope {
205 return s.delegate.(network.PeerScope)
206 }
207
208 func (s *loggingScope) SetPeer(p peer.ID) error {
209 err := s.delegate.(network.ConnManagementScope).SetPeer(p)
210 s.countErrs(err)
211 return err
212 }
213
214 func (s *loggingScope) ProtocolScope() network.ProtocolScope {
215 return s.delegate.(network.ProtocolScope)
216 }
217
218 func (s *loggingScope) SetProtocol(proto protocol.ID) error {
219 err := s.delegate.(network.StreamManagementScope).SetProtocol(proto)
220 s.countErrs(err)
221 return err
222 }
223
224 func (s *loggingScope) ServiceScope() network.ServiceScope {
225 return s.delegate.(network.ServiceScope)
226 }
227
228 func (s *loggingScope) SetService(srv string) error {
229 err := s.delegate.(network.StreamManagementScope).SetService(srv)
230 s.countErrs(err)
231 return err
232 }
233
234 func (s *loggingScope) Limit() rcmgr.Limit {
235 return s.delegate.(rcmgr.ResourceScopeLimiter).Limit()
236 }
237
238 func (s *loggingScope) SetLimit(limit rcmgr.Limit) {
239 s.delegate.(rcmgr.ResourceScopeLimiter).SetLimit(limit)
240 }