| 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 | } |