"github.com/anacrolix/missinggo/v2/pubsub"
"github.com/anacrolix/multiless"
"github.com/anacrolix/sync"
- "github.com/davecgh/go-spew/spew"
"github.com/pion/datachannel"
"golang.org/x/exp/maps"
"github.com/anacrolix/torrent/bencode"
"github.com/anacrolix/torrent/common"
+ "github.com/anacrolix/torrent/internal/check"
+ "github.com/anacrolix/torrent/internal/nestedmaps"
"github.com/anacrolix/torrent/metainfo"
pp "github.com/anacrolix/torrent/peer_protocol"
utHolepunch "github.com/anacrolix/torrent/peer_protocol/ut-holepunch"
maxEstablishedConns int
// Set of addrs to which we're attempting to connect. Connections are
// half-open until all handshakes are completed.
- halfOpen map[string]PeerInfo
+ halfOpen map[string]map[outgoingConnAttemptKey]*PeerInfo
// The final ess is not silent here as it's in the plural.
utHolepunchRendezvous map[netip.AddrPort]*utHolepunchRendezvous
requestIndexes []RequestIndex
}
+type outgoingConnAttemptKey = *PeerInfo
+
func (t *Torrent) length() int64 {
return t._length.Value
}
})
// Add half-open peers to the list
- for _, peer := range t.halfOpen {
- ks = append(ks, peer)
+ for _, attempts := range t.halfOpen {
+ for _, peer := range attempts {
+ ks = append(ks, *peer)
+ }
}
// Add active peers to the list
fmt.Fprintf(w, "DHT Announces: %d\n", t.numDHTAnnounces)
- spew.NewDefaultConfig()
- spew.Fdump(w, t.statsLocked())
+ dumpStats(w, t.statsLocked())
fmt.Fprintf(w, "webseeds:\n")
t.writePeerStatuses(w, maps.Values(t.webSeeds))
return ml.Less()
})
- fmt.Fprintf(w, "peer conns:\n")
+ fmt.Fprintf(w, "%v peer conns:\n", len(peerConns))
t.writePeerStatuses(w, g.SliceMap(peerConns, func(pc *PeerConn) *Peer {
return &pc.Peer
}))
return
}
p := t.peers.PopMax()
- t.initiateConn(p, false, false)
+ t.initiateConn(p, false, false, false, false)
initiated++
}
return
}
func (t *Torrent) assertPendingRequests() {
- if !check {
+ if !check.Enabled {
return
}
// var actual pendingRequests
}
}
+func (t *Torrent) connectingToPeerAddr(addrStr string) bool {
+ return len(t.halfOpen[addrStr]) != 0
+}
+
+func (t *Torrent) hasPeerConnForAddr(x PeerRemoteAddr) bool {
+ addrStr := x.String()
+ for c := range t.conns {
+ ra := c.RemoteAddr
+ if ra.String() == addrStr {
+ return true
+ }
+ }
+ return false
+}
+
+func (t *Torrent) getHalfOpenPath(
+ addrStr string,
+ attemptKey outgoingConnAttemptKey,
+) nestedmaps.Path[*PeerInfo] {
+ return nestedmaps.Next(nestedmaps.Next(nestedmaps.Begin(&t.halfOpen), addrStr), attemptKey)
+}
+
+func (t *Torrent) addHalfOpen(addrStr string, attemptKey *PeerInfo) {
+ path := t.getHalfOpenPath(addrStr, attemptKey)
+ if path.Exists() {
+ panic("should be unique")
+ }
+ path.Set(attemptKey)
+ t.cl.numHalfOpen++
+}
+
// Start the process of connecting to the given peer for the given torrent if appropriate. I'm not
// sure all the PeerInfo fields are being used.
func (t *Torrent) initiateConn(
peer PeerInfo,
requireRendezvous bool,
skipHolepunchRendezvous bool,
+ ignoreLimits bool,
+ receivedHolepunchConnect bool,
) {
if peer.Id == t.cl.peerID {
return
return
}
addr := peer.Addr
- if t.addrActive(addr.String()) {
+ addrStr := addr.String()
+ if !ignoreLimits {
+ if t.connectingToPeerAddr(addrStr) {
+ return
+ }
+ }
+ if t.hasPeerConnForAddr(addr) {
return
}
- t.cl.numHalfOpen++
- t.halfOpen[addr.String()] = peer
- go t.cl.outgoingConnection(outgoingConnOpts{
- t: t,
- addr: peer.Addr,
- requireRendezvous: requireRendezvous,
- skipHolepunchRendezvous: skipHolepunchRendezvous,
- }, peer.Source, peer.Trusted)
+ attemptKey := &peer
+ t.addHalfOpen(addrStr, attemptKey)
+ go t.cl.outgoingConnection(
+ outgoingConnOpts{
+ t: t,
+ addr: peer.Addr,
+ requireRendezvous: requireRendezvous,
+ skipHolepunchRendezvous: skipHolepunchRendezvous,
+ receivedHolepunchConnect: receivedHolepunchConnect,
+ },
+ peer.Source,
+ peer.Trusted,
+ attemptKey,
+ )
}
// Adds a trusted, pending peer for each of the given Client's addresses. Typically used in tests to
// connection.
func (t *Torrent) allStats(f func(*ConnStats)) {
f(&t.stats)
- f(&t.cl.stats)
+ f(&t.cl.connStats)
}
func (t *Torrent) hashingPiece(i pieceIndex) bool {
t.initiateConn(PeerInfo{
Addr: msg.AddrPort,
Source: PeerSourceUtHolepunch,
- }, false, true)
+ }, false, true, true, true)
}
return nil
case utHolepunch.Error:
}
return
}
+
+func (t *Torrent) numHalfOpenAttempts() (num int) {
+ for _, attempts := range t.halfOpen {
+ num += len(attempts)
+ }
+ return
+}