"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
}
t.conns[c] = struct{}{}
t.cl.event.Broadcast()
+ // We'll never receive the "p" extended handshake parameter.
if !t.cl.config.DisablePEX && !c.PeerExtensionBytes.SupportsExtended() {
- t.pex.Add(c) // as no further extended handshake expected
+ t.pex.Add(c)
}
return nil
}
}
}
+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 {
return nil
}
-func (t *Torrent) peerConnsWithRemoteAddrPort(addrPort netip.AddrPort) (ret []*PeerConn) {
+func (t *Torrent) peerConnsWithDialAddrPort(target netip.AddrPort) (ret []*PeerConn) {
for pc := range t.conns {
- addr := pc.remoteAddrPort()
- if !(addr.Ok && addr.Value == addrPort) {
+ dialAddr, err := pc.remoteDialAddrPort()
+ if err != nil {
+ continue
+ }
+ if dialAddr != target {
continue
}
ret = append(ret, pc)
case utHolepunch.Rendezvous:
t.logger.Printf("got holepunch rendezvous request for %v from %p", msg.AddrPort, sender)
sendMsg := sendUtHolepunchMsg
- targets := t.peerConnsWithRemoteAddrPort(msg.AddrPort)
+ senderAddrPort, err := sender.remoteDialAddrPort()
+ if err != nil {
+ sender.logger.Levelf(
+ log.Warning,
+ "error getting ut_holepunch rendezvous sender's dial address: %v",
+ err,
+ )
+ // There's no better error code. The sender's address itself is invalid. I don't see
+ // this error message being appropriate anywhere else anyway.
+ sendMsg(sender, utHolepunch.Error, msg.AddrPort, utHolepunch.NoSuchPeer)
+ }
+ targets := t.peerConnsWithDialAddrPort(msg.AddrPort)
if len(targets) == 0 {
sendMsg(sender, utHolepunch.Error, msg.AddrPort, utHolepunch.NotConnected)
return nil
continue
}
sendMsg(sender, utHolepunch.Connect, msg.AddrPort, 0)
- sendMsg(pc, utHolepunch.Connect, sender.remoteAddrPort().Unwrap(), 0)
+ sendMsg(pc, utHolepunch.Connect, senderAddrPort, 0)
}
return nil
case utHolepunch.Connect:
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
+}