peerconn.go | 18 ++++++++++++------ request-strategy.go | 4 ++++ diff --git a/peerconn.go b/peerconn.go index 71ac5f3d252a230d4b7bbed38ae9033e41d30e52..f3680e1832472b6e893c2784a4ddcc621bc42dd1 100644 --- a/peerconn.go +++ b/peerconn.go @@ -47,6 +47,10 @@ type PeerRemoteAddr interface { String() string } +// Since we have to store all the requests in memory, we can't reasonably exceed what would be +// indexable with the memory space available. +type maxRequests = int + type Peer struct { // First to ensure 64-bit alignment for atomics. See #262. _stats ConnStats @@ -83,9 +87,10 @@ lastStartedExpectingToReceiveChunks time.Time cumulativeExpectedToReceiveChunks time.Duration _chunksReceivedWhileExpecting int64 - choking bool - requests map[Request]struct{} - requestsLowWater int + choking bool + requests map[Request]struct{} + piecesReceivedSinceLastRequestUpdate maxRequests + maxPiecesReceivedBetweenRequestUpdates maxRequests // Chunks that we might reasonably expect to receive from the peer. Due to // latency, buffering, and implementation differences, we may receive // chunks that are no longer in the set of requests actually want. @@ -114,7 +119,7 @@ // Pieces we've accepted chunks for from the peer. peerTouchedPieces map[pieceIndex]struct{} peerAllowedFast bitmap.Bitmap - PeerMaxRequests int // Maximum pending requests the peer allows. + PeerMaxRequests maxRequests // Maximum pending requests the peer allows. PeerExtensionIDs map[pp.ExtensionName]pp.ExtensionNumber PeerClientName string @@ -470,8 +475,8 @@ return index < len(cn.metadataRequests) && cn.metadataRequests[index] } // The actual value to use as the maximum outbound requests. -func (cn *Peer) nominalMaxRequests() (ret int) { - return int(clamp(1, int64(cn.PeerMaxRequests), 64)) +func (cn *Peer) nominalMaxRequests() (ret maxRequests) { + return int(clamp(1, 2*int64(cn.maxPiecesReceivedBetweenRequestUpdates), int64(cn.PeerMaxRequests))) } func (cn *Peer) totalExpectingTime() (ret time.Duration) { @@ -1358,6 +1363,7 @@ c.allStats(add(1, func(cs *ConnStats) *Count { return &cs.ChunksReadUseful })) c.allStats(add(int64(len(msg.Piece)), func(cs *ConnStats) *Count { return &cs.BytesReadUsefulData })) if deletedRequest { + c.piecesReceivedSinceLastRequestUpdate++ c.allStats(add(int64(len(msg.Piece)), func(cs *ConnStats) *Count { return &cs.BytesReadUsefulIntendedData })) } for _, f := range c.t.cl.config.Callbacks.ReceivedUsefulData { diff --git a/request-strategy.go b/request-strategy.go index 7c7660bc4cc83d705a8c38541c5ad835fff2a0d7..f0470007ccf91d58567c2b54b62f3862b0a96751 100644 --- a/request-strategy.go +++ b/request-strategy.go @@ -53,6 +53,10 @@ t.iterPeers(func(p *Peer) { if p.closed.IsSet() { return } + if p.piecesReceivedSinceLastRequestUpdate > p.maxPiecesReceivedBetweenRequestUpdates { + p.maxPiecesReceivedBetweenRequestUpdates = p.piecesReceivedSinceLastRequestUpdate + } + p.piecesReceivedSinceLastRequestUpdate = 0 rst.Peers = append(rst.Peers, request_strategy.Peer{ HasPiece: p.peerHasPiece, MaxRequests: p.nominalMaxRequests(),