peerconn.go | 42 ++++++++++++++++++++++++++---------------- torrent.go | 23 ++++++++++++++++++++++- web_seed.go | 33 +++++++++++++++++++++++++++++++++ diff --git a/peerconn.go b/peerconn.go index f7a995919e89708379d6c85039be39ef8119d2fc..9ec40a28d2305314f9ca92e49bcec00b6c4713ce 100644 --- a/peerconn.go +++ b/peerconn.go @@ -41,6 +41,8 @@ writeInterested(interested bool) bool cancel(request) bool request(request) bool connectionFlags() string + close() + postCancel(request) drop() } @@ -335,16 +337,20 @@ }(), ) } -func (cn *PeerConn) close() { +func (cn *peer) close() { if !cn.closed.Set() { return } + cn.discardPieceInclination() + cn._pieceRequestOrder.Clear() + cn.peerImpl.close() +} + +func (cn *PeerConn) close() { if cn.pex.IsEnabled() { cn.pex.Close() } cn.tickleWriter() - cn.discardPieceInclination() - cn._pieceRequestOrder.Clear() if cn.conn != nil { cn.conn.Close() } @@ -365,6 +371,15 @@ // Last I checked only Piece messages affect stats, and we don't post // those. cn.wroteMsg(&msg) cn.tickleWriter() +} + +func (cn *PeerConn) write(msg pp.Message) bool { + cn.wroteMsg(&msg) + cn.writeBuffer.Write(msg.MustMarshalBinary()) + torrent.Add(fmt.Sprintf("messages filled of type %s", msg.Type.String()), 1) + // 64KiB, but temporarily less to work around an issue with WebRTC. TODO: Update + // when https://github.com/pion/datachannel/issues/59 is fixed. + return cn.writeBuffer.Len() < 1<<15 } func (cn *PeerConn) requestMetadataPiece(index int) { @@ -606,15 +621,6 @@ } cn.upload(cn.write) } -func (cn *PeerConn) write(msg pp.Message) bool { - cn.wroteMsg(&msg) - cn.writeBuffer.Write(msg.MustMarshalBinary()) - torrent.Add(fmt.Sprintf("messages filled of type %s", msg.Type.String()), 1) - // 64KiB, but temporarily less to work around an issue with WebRTC. TODO: Update - // when https://github.com/pion/datachannel/issues/59 is fixed. - return cn.writeBuffer.Len() < 1<<15 -} - // Routine that writes to the peer. Some of what to write is buffered by // activity elsewhere in the Client, and some is determined locally when the // connection is writable. @@ -803,7 +809,7 @@ } return cn.pieceInclination } -func (cn *PeerConn) discardPieceInclination() { +func (cn *peer) discardPieceInclination() { if cn.pieceInclination == nil { return } @@ -1475,7 +1481,7 @@ }) return true } -func (c *PeerConn) deleteAllRequests() { +func (c *peer) deleteAllRequests() { for r := range c.requests { c.deleteRequest(r) } @@ -1491,12 +1497,16 @@ func (c *PeerConn) tickleWriter() { c.writerCond.Broadcast() } -func (c *PeerConn) postCancel(r request) bool { +func (c *peer) postCancel(r request) bool { if !c.deleteRequest(r) { return false } - c.post(makeCancelMessage(r)) + c.peerImpl.postCancel(r) return true +} + +func (c *PeerConn) postCancel(r request) { + c.post(makeCancelMessage(r)) } func (c *PeerConn) sendChunk(r request, msg func(pp.Message) bool) (more bool, err error) { diff --git a/torrent.go b/torrent.go index 7b1d1db31494987f39149f49af582ee51f1b3d15..3f1b8f039c333c405a684820b0c0a742479ea087 100644 --- a/torrent.go +++ b/torrent.go @@ -77,6 +77,8 @@ // The info dict. nil if we don't have it (yet). info *metainfo.Info files *[]*File + webSeeds map[string]*peer + // Active peer connections, running message stream loops. TODO: Make this // open (not-closed) connections only. conns map[*PeerConn]struct{} @@ -1226,9 +1228,16 @@ t.pex.Drop(c) } torrent.Add("deleted connections", 1) c.deleteAllRequests() - if len(t.conns) == 0 { + if t.numActivePeers() == 0 { t.assertNoPendingRequests() } + return +} + +func (t *Torrent) numActivePeers() (num int) { + t.iterPeers(func(*peer) { + num++ + }) return } @@ -1978,6 +1987,18 @@ func (t *Torrent) iterPeers(f func(*peer)) { for pc := range t.conns { f(&pc.peer) + } + for _, ws := range t.webSeeds { + f(ws) + } +} + +func (t *Torrent) addWebSeed(url string) { + if _, ok := t.webSeeds[url]; ok { + return + } + t.webSeeds[url] = &peer{ + peerImpl: &webSeed{}, } } diff --git a/web_seed.go b/web_seed.go new file mode 100644 index 0000000000000000000000000000000000000000..2242f1c8757e1a95363490b73eddf11e1ad8d60f --- /dev/null +++ b/web_seed.go @@ -0,0 +1,33 @@ +package torrent + +import ( + "net/http" +) + +type webSeed struct { + peer *peer + httpClient *http.Client +} + +func (ws *webSeed) writeInterested(interested bool) bool { + return true +} + +func (ws *webSeed) cancel(r request) bool { + panic("implement me") +} + +func (ws *webSeed) request(r request) bool { + panic("implement me") +} + +func (ws *webSeed) connectionFlags() string { + return "WS" +} + +func (ws *webSeed) drop() { +} + +func (ws *webSeed) updateRequests() { + ws.peer.doRequestState() +}