2021-05-13 07:56:58 +08:00
|
|
|
package request_strategy
|
|
|
|
|
|
|
|
import (
|
2021-09-18 18:34:14 +08:00
|
|
|
"bytes"
|
2021-11-30 12:18:38 +08:00
|
|
|
"expvar"
|
2021-05-13 07:56:58 +08:00
|
|
|
|
|
|
|
"github.com/anacrolix/multiless"
|
2021-11-30 12:18:38 +08:00
|
|
|
"github.com/anacrolix/torrent/metainfo"
|
2021-12-01 11:38:47 +08:00
|
|
|
"github.com/google/btree"
|
2021-05-21 12:02:45 +08:00
|
|
|
|
2021-05-13 07:56:58 +08:00
|
|
|
"github.com/anacrolix/torrent/types"
|
|
|
|
)
|
|
|
|
|
|
|
|
type (
|
2021-09-19 13:16:37 +08:00
|
|
|
RequestIndex = uint32
|
|
|
|
ChunkIndex = uint32
|
2021-05-13 07:56:58 +08:00
|
|
|
Request = types.Request
|
|
|
|
pieceIndex = types.PieceIndex
|
|
|
|
piecePriority = types.PiecePriority
|
2021-05-13 09:26:22 +08:00
|
|
|
// This can be made into a type-param later, will be great for testing.
|
|
|
|
ChunkSpec = types.ChunkSpec
|
2021-05-13 07:56:58 +08:00
|
|
|
)
|
|
|
|
|
2021-12-01 11:38:47 +08:00
|
|
|
type pieceOrderInput struct {
|
|
|
|
PieceRequestOrderState
|
|
|
|
PieceRequestOrderKey
|
|
|
|
}
|
|
|
|
|
|
|
|
func pieceOrderLess(i, j pieceOrderInput) multiless.Computation {
|
2021-11-30 18:31:32 +08:00
|
|
|
return multiless.New().Int(
|
|
|
|
int(j.Priority), int(i.Priority),
|
|
|
|
).Bool(
|
|
|
|
j.Partial, i.Partial,
|
|
|
|
).Int64(
|
|
|
|
i.Availability, j.Availability,
|
|
|
|
).Int(
|
2021-12-01 11:38:47 +08:00
|
|
|
i.Index, j.Index,
|
2021-11-30 18:31:32 +08:00
|
|
|
).Lazy(func() multiless.Computation {
|
|
|
|
return multiless.New().Cmp(bytes.Compare(
|
2021-12-01 11:38:47 +08:00
|
|
|
i.InfoHash[:],
|
|
|
|
j.InfoHash[:],
|
2021-11-30 18:31:32 +08:00
|
|
|
))
|
2021-12-01 11:38:47 +08:00
|
|
|
})
|
2021-05-13 07:56:58 +08:00
|
|
|
}
|
|
|
|
|
|
|
|
type requestsPeer struct {
|
2021-05-13 09:26:22 +08:00
|
|
|
Peer
|
2021-05-13 07:56:58 +08:00
|
|
|
nextState PeerNextRequestState
|
|
|
|
requestablePiecesRemaining int
|
|
|
|
}
|
|
|
|
|
|
|
|
func (rp *requestsPeer) canFitRequest() bool {
|
2021-09-19 13:16:37 +08:00
|
|
|
return int(rp.nextState.Requests.GetCardinality()) < rp.MaxRequests
|
2021-05-13 07:56:58 +08:00
|
|
|
}
|
|
|
|
|
2021-09-19 13:16:37 +08:00
|
|
|
func (rp *requestsPeer) addNextRequest(r RequestIndex) {
|
|
|
|
if !rp.nextState.Requests.CheckedAdd(r) {
|
2021-05-13 18:56:12 +08:00
|
|
|
panic("should only add once")
|
2021-05-13 07:56:58 +08:00
|
|
|
}
|
|
|
|
}
|
|
|
|
|
|
|
|
type peersForPieceRequests struct {
|
|
|
|
requestsInPiece int
|
|
|
|
*requestsPeer
|
|
|
|
}
|
|
|
|
|
2021-09-19 13:16:37 +08:00
|
|
|
func (me *peersForPieceRequests) addNextRequest(r RequestIndex) {
|
2021-05-13 18:56:12 +08:00
|
|
|
me.requestsPeer.addNextRequest(r)
|
|
|
|
me.requestsInPiece++
|
2021-05-13 07:56:58 +08:00
|
|
|
}
|
|
|
|
|
2021-11-30 12:18:38 +08:00
|
|
|
var packageExpvarMap = expvar.NewMap("request-strategy")
|
|
|
|
|
2021-09-18 16:57:50 +08:00
|
|
|
// Calls f with requestable pieces in order.
|
2021-12-01 16:21:25 +08:00
|
|
|
func GetRequestablePieces(input Input, pro *PieceRequestOrder, f func(ih metainfo.Hash, pieceIndex int)) {
|
2021-05-13 07:56:58 +08:00
|
|
|
// Storage capacity left for this run, keyed by the storage capacity pointer on the storage
|
2021-09-15 08:30:37 +08:00
|
|
|
// TorrentImpl. A nil value means no capacity limit.
|
2021-11-29 10:07:18 +08:00
|
|
|
var storageLeft *int64
|
2021-12-01 16:21:25 +08:00
|
|
|
if cap, ok := input.Capacity(); ok {
|
|
|
|
storageLeft = &cap
|
2021-11-29 10:07:18 +08:00
|
|
|
}
|
2021-05-14 11:40:09 +08:00
|
|
|
var allTorrentsUnverifiedBytes int64
|
2021-12-01 11:38:47 +08:00
|
|
|
pro.tree.Ascend(func(i btree.Item) bool {
|
|
|
|
_i := i.(pieceRequestOrderItem)
|
2021-12-01 16:21:25 +08:00
|
|
|
ih := _i.key.InfoHash
|
|
|
|
var t Torrent = input.Torrent(ih)
|
|
|
|
var piece Piece = t.Piece(_i.key.Index)
|
|
|
|
pieceLength := t.PieceLength()
|
|
|
|
if storageLeft != nil {
|
|
|
|
if *storageLeft < pieceLength {
|
2021-12-01 11:38:47 +08:00
|
|
|
return true
|
2021-05-13 07:56:58 +08:00
|
|
|
}
|
2021-12-01 16:21:25 +08:00
|
|
|
*storageLeft -= pieceLength
|
2021-05-14 09:50:41 +08:00
|
|
|
}
|
2021-12-01 16:21:25 +08:00
|
|
|
if !piece.Request() || piece.NumPendingChunks() == 0 {
|
2021-05-14 11:40:09 +08:00
|
|
|
// TODO: Clarify exactly what is verified. Stuff that's being hashed should be
|
|
|
|
// considered unverified and hold up further requests.
|
2021-12-01 11:38:47 +08:00
|
|
|
return true
|
2021-05-13 07:56:58 +08:00
|
|
|
}
|
2021-12-01 16:21:25 +08:00
|
|
|
if input.MaxUnverifiedBytes() != 0 && allTorrentsUnverifiedBytes+pieceLength > input.MaxUnverifiedBytes() {
|
2021-12-01 11:38:47 +08:00
|
|
|
return true
|
2021-05-14 11:40:09 +08:00
|
|
|
}
|
2021-12-01 16:21:25 +08:00
|
|
|
allTorrentsUnverifiedBytes += pieceLength
|
|
|
|
f(ih, _i.key.Index)
|
2021-12-01 11:38:47 +08:00
|
|
|
return true
|
|
|
|
})
|
2021-05-14 11:06:12 +08:00
|
|
|
return
|
|
|
|
}
|
|
|
|
|
2021-12-01 16:21:25 +08:00
|
|
|
type Input interface {
|
|
|
|
Torrent(metainfo.Hash) Torrent
|
|
|
|
// Storage capacity, shared among all Torrents with the same storage.TorrentCapacity pointer in
|
|
|
|
// their storage.Torrent references.
|
|
|
|
Capacity() (cap int64, capped bool)
|
2021-11-29 10:07:18 +08:00
|
|
|
// Across all the Torrents. This might be partitioned by storage capacity key now.
|
2021-12-01 16:21:25 +08:00
|
|
|
MaxUnverifiedBytes() int64
|
2021-05-13 16:35:49 +08:00
|
|
|
}
|