client.go | 37 +++++-------------------------------- client_test.go | 15 ++++++++------- data.go | 15 +++++++++++++++ data/blob/store.go | 7 ++++--- data/data.go | 21 --------------------- data/file/file.go | 18 +++++++++++++++--- data/mmap/mmap.go | 26 +++++++++++++++++++++++--- fs/torrentfs_test.go | 6 ++---- stateless.go | 21 --------------------- torrent.go | 24 +++++------------------- diff --git a/client.go b/client.go index 898ee23a36cb9c397e89a33726dd1bd946d9143b..d28357269ec8b4cdfaf8cca73a83dc882f35120c 100644 --- a/client.go +++ b/client.go @@ -32,7 +32,6 @@ "github.com/bradfitz/iter" "github.com/edsrzf/mmap-go" "github.com/anacrolix/torrent/bencode" - "github.com/anacrolix/torrent/data" filePkg "github.com/anacrolix/torrent/data/file" "github.com/anacrolix/torrent/dht" "github.com/anacrolix/torrent/internal/pieceordering" @@ -253,40 +252,14 @@ fmt.Fprintln(w) } } -// A Data that implements this has a streaming interface that should be -// preferred over ReadAt. For example, the data is stored in blocks on the -// network and have a fixed cost to open. -type SectionOpener interface { - // Open a ReadCloser at the given offset into torrent data. n is how many - // bytes we intend to read. - OpenSection(off, n int64) (io.ReadCloser, error) -} - -func dataReadAt(d data.Data, b []byte, off int64) (n int, err error) { +func dataReadAt(d Data, b []byte, off int64) (n int, err error) { // defer func() { // if err == io.ErrUnexpectedEOF && n != 0 { // err = nil // } // }() // log.Println("data read at", len(b), off) -again: - if ra, ok := d.(io.ReaderAt); ok { - return ra.ReadAt(b, off) - } - if so, ok := d.(SectionOpener); ok { - var rc io.ReadCloser - rc, err = so.OpenSection(off, int64(len(b))) - if err != nil { - return - } - defer rc.Close() - return io.ReadFull(rc, b) - } - if dp, ok := super(d); ok { - d = dp.(data.Data) - goto again - } - panic(fmt.Sprintf("can't read from %T", d)) + return d.ReadAt(b, off) } // Calculates the number of pieces to set to Readahead priority, after the @@ -474,7 +447,7 @@ }() cl = &Client{ halfOpenLimit: socketsPerTorrent, config: *cfg, - torrentDataOpener: func(md *metainfo.Info) data.Data { + torrentDataOpener: func(md *metainfo.Info) Data { return filePkg.TorrentData(md, cfg.DataDir) }, dopplegangerAddrs: make(map[string]struct{}), @@ -1934,7 +1907,7 @@ } } // Storage cannot be changed once it's set. -func (cl *Client) setStorage(t *torrent, td data.Data) (err error) { +func (cl *Client) setStorage(t *torrent, td Data) (err error) { err = t.setStorage(td) cl.event.Broadcast() if err != nil { @@ -1944,7 +1917,7 @@ cl.startTorrent(t) return } -type TorrentDataOpener func(*metainfo.Info) data.Data +type TorrentDataOpener func(*metainfo.Info) Data func (cl *Client) setMetaData(t *torrent, md *metainfo.Info, bytes []byte) (err error) { err = t.setMetadata(md, bytes, &cl.mu) diff --git a/client_test.go b/client_test.go index fb797b2bb9607e1f30602d028a6f202b14495bac..49d0965fd77f827bdeba84fd23c30a712bd8adce 100644 --- a/client_test.go +++ b/client_test.go @@ -22,7 +22,6 @@ "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" "github.com/anacrolix/torrent/bencode" - "github.com/anacrolix/torrent/data" "github.com/anacrolix/torrent/data/blob" "github.com/anacrolix/torrent/dht" "github.com/anacrolix/torrent/internal/testutil" @@ -267,7 +266,10 @@ defer os.RemoveAll(leecherDataDir) // cfg.TorrentDataOpener = func(info *metainfo.Info) (data.Data, error) { // return blob.TorrentData(info, leecherDataDir), nil // } - cfg.TorrentDataOpener = blob.NewStore(leecherDataDir).OpenTorrent + blobStore := blob.NewStore(leecherDataDir) + cfg.TorrentDataOpener = func(info *metainfo.Info) Data { + return blobStore.OpenTorrent(info) + } leecher, _ := NewClient(&cfg) defer leecher.Close() leecherGreeting, _, _ := leecher.AddTorrentSpec(func() (ret *TorrentSpec) { @@ -402,8 +404,9 @@ assert.EqualValues(t, T.Trackers[0][0].URL(), "http://a") assert.EqualValues(t, T.Trackers[1][0].URL(), "udp://b") } -type badData struct { -} +type badData struct{} + +func (me badData) Close() {} func (me badData) WriteAt(b []byte, off int64) (int, error) { return 0, nil @@ -430,12 +433,10 @@ n = copy(b, []byte("hello")[off:]) return } -var _ StatefulData = badData{} - // We read from a piece which is marked completed, but is missing data. func TestCompletedPieceWrongSize(t *testing.T) { cfg := TestingConfig - cfg.TorrentDataOpener = func(*metainfo.Info) data.Data { + cfg.TorrentDataOpener = func(*metainfo.Info) Data { return badData{} } cl, _ := NewClient(&cfg) diff --git a/data.go b/data.go new file mode 100644 index 0000000000000000000000000000000000000000..803d8476bbb0da4c1fc9bd38fdc462de7e9301af --- /dev/null +++ b/data.go @@ -0,0 +1,15 @@ +package torrent + +import "io" + +// Represents data storage for a Torrent. +type Data interface { + ReadAt(p []byte, off int64) (n int, err error) + Close() + WriteAt(p []byte, off int64) (n int, err error) + WriteSectionTo(w io.Writer, off, n int64) (written int64, err error) + // We believe the piece data will pass a hash check. + PieceCompleted(index int) error + // Returns true if the piece is complete. + PieceComplete(index int) bool +} diff --git a/data/blob/store.go b/data/blob/store.go index 96eda1feb8cf121d405bd064790d074be22b4d9c..385d2ee5a1d9ca91e491652aaeef7434a130c49e 100644 --- a/data/blob/store.go +++ b/data/blob/store.go @@ -1,3 +1,4 @@ +// Implements torrent data storage as per-piece files. package blob import ( @@ -14,8 +15,8 @@ "sync" "time" "github.com/anacrolix/missinggo" + "github.com/anacrolix/torrent" - dataPkg "github.com/anacrolix/torrent/data" "github.com/anacrolix/torrent/metainfo" ) @@ -32,7 +33,7 @@ mu sync.Mutex completed map[[20]byte]struct{} } -func (me *store) OpenTorrent(info *metainfo.Info) dataPkg.Data { +func (me *store) OpenTorrent(info *metainfo.Info) torrent.Data { return &data{info, me} } @@ -44,7 +45,7 @@ s.capacity = bytes } } -func NewStore(baseDir string, opt ...StoreOption) dataPkg.Store { +func NewStore(baseDir string, opt ...StoreOption) *store { s := &store{baseDir, -1, sync.Mutex{}, nil} for _, o := range opt { o(s) diff --git a/data/data.go b/data/data.go deleted file mode 100644 index f783279ddf0a654b92f521ff54dafc7b42b8b732..0000000000000000000000000000000000000000 --- a/data/data.go +++ /dev/null @@ -1,21 +0,0 @@ -package data - -import ( - "io" - - "github.com/anacrolix/torrent/metainfo" -) - -type Store interface { - OpenTorrent(*metainfo.Info) Data -} - -// Represents data storage for a Torrent. Additional optional interfaces to -// implement are io.Closer, io.ReaderAt, StatefulData, and SectionOpener. -type Data interface { - // OpenSection(off, n int64) (io.ReadCloser, error) - // ReadAt(p []byte, off int64) (n int, err error) - // Close() - WriteAt(p []byte, off int64) (n int, err error) - WriteSectionTo(w io.Writer, off, n int64) (written int64, err error) -} diff --git a/data/file/file.go b/data/file/file.go index 0c30bc50c14aa7c289eb871a6c9aa10df4c6faf6..bf77db2adf5ffcbdf800425bfe956f6ee2e2c87b 100644 --- a/data/file/file.go +++ b/data/file/file.go @@ -9,12 +9,24 @@ "github.com/anacrolix/torrent/metainfo" ) type data struct { - info *metainfo.Info - loc string + info *metainfo.Info + loc string + completed []bool } func TorrentData(md *metainfo.Info, location string) data { - return data{md, location} + return data{md, location, make([]bool, md.NumPieces())} +} + +func (me data) Close() {} + +func (me data) PieceComplete(piece int) bool { + return me.completed[piece] +} + +func (me data) PieceCompleted(piece int) error { + me.completed[piece] = true + return nil } func (me data) ReadAt(p []byte, off int64) (n int, err error) { diff --git a/data/mmap/mmap.go b/data/mmap/mmap.go index 7a0324b3655d4fbbe03752c094215af9696e6d8b..ec4236d64c817478a9f1e781341ce80e0de20bf9 100644 --- a/data/mmap/mmap.go +++ b/data/mmap/mmap.go @@ -11,12 +11,28 @@ "github.com/anacrolix/torrent/metainfo" "github.com/anacrolix/torrent/mmap_span" ) -func TorrentData(md *metainfo.Info, location string) (mms *mmap_span.MMapSpan, err error) { - mms = &mmap_span.MMapSpan{} +type torrentData struct { + // Supports non-torrent specific data operations for the torrent.Data + // interface. + mmap_span.MMapSpan + + completed []bool +} + +func (me *torrentData) PieceComplete(piece int) bool { + return me.completed[piece] +} + +func (me *torrentData) PieceCompleted(piece int) error { + me.completed[piece] = true + return nil +} + +func TorrentData(md *metainfo.Info, location string) (ret *torrentData, err error) { + var mms mmap_span.MMapSpan defer func() { if err != nil { mms.Close() - mms = nil } }() for _, miFile := range md.UpvertedFiles() { @@ -62,6 +78,10 @@ }() if err != nil { return } + } + ret = &torrentData{ + MMapSpan: mms, + completed: make([]bool, md.NumPieces()), } return } diff --git a/fs/torrentfs_test.go b/fs/torrentfs_test.go index 8ae89e1c30bd33237942751f0b8767f91315fc08..559e79886e7df3dd3710ffedfb934ab78ef219e0 100644 --- a/fs/torrentfs_test.go +++ b/fs/torrentfs_test.go @@ -15,16 +15,14 @@ "strings" "testing" "time" - _ "github.com/anacrolix/envpprof" - "bazil.org/fuse" fusefs "bazil.org/fuse/fs" + _ "github.com/anacrolix/envpprof" "github.com/anacrolix/missinggo" "github.com/stretchr/testify/assert" netContext "golang.org/x/net/context" "github.com/anacrolix/torrent" - "github.com/anacrolix/torrent/data" "github.com/anacrolix/torrent/data/mmap" "github.com/anacrolix/torrent/internal/testutil" "github.com/anacrolix/torrent/metainfo" @@ -205,7 +203,7 @@ DisableTCP: true, NoDefaultBlocklist: true, - TorrentDataOpener: func(info *metainfo.Info) data.Data { + TorrentDataOpener: func(info *metainfo.Info) Data { ret, _ := mmap.TorrentData(info, filepath.Join(layout.BaseDir, "download")) return ret }, diff --git a/stateless.go b/stateless.go deleted file mode 100644 index 99fc1e071144c5e6a2e2aaff4974664b6b13fc50..0000000000000000000000000000000000000000 --- a/stateless.go +++ /dev/null @@ -1,21 +0,0 @@ -package torrent - -import "github.com/anacrolix/torrent/data" - -type statelessDataWrapper struct { - data.Data - complete []bool -} - -func (me *statelessDataWrapper) PieceComplete(piece int) bool { - return me.complete[piece] -} - -func (me *statelessDataWrapper) PieceCompleted(piece int) error { - me.complete[piece] = true - return nil -} - -func (me *statelessDataWrapper) Super() interface{} { - return me.Data -} diff --git a/torrent.go b/torrent.go index 82d97a49051def3514e0de8c219be56a80f4ebb9..cb130007c774ef8dd454c0cdec6e14a917097dd2 100644 --- a/torrent.go +++ b/torrent.go @@ -17,7 +17,6 @@ "github.com/anacrolix/missinggo/pubsub" "github.com/bradfitz/iter" "github.com/anacrolix/torrent/bencode" - "github.com/anacrolix/torrent/data" "github.com/anacrolix/torrent/metainfo" pp "github.com/anacrolix/torrent/peer_protocol" "github.com/anacrolix/torrent/tracker" @@ -45,15 +44,6 @@ IPBytes string Port int } -// Data maintains per-piece persistent state. -type StatefulData interface { - data.Data - // We believe the piece data will pass a hash check. - PieceCompleted(index int) error - // Returns true if the piece is complete. - PieceComplete(index int) bool -} - // Is not aware of Client. Maintains state of torrent for with-in a Client. type torrent struct { stateMu sync.Mutex @@ -75,7 +65,7 @@ // Total length of the torrent in bytes. Stored because it's not O(1) to // get this from the info dict. length int64 - data StatefulData + data Data // The info dict. Nil if we don't have it (yet). Info *metainfo.Info @@ -272,15 +262,11 @@ } return } -func (t *torrent) setStorage(td data.Data) (err error) { - if c, ok := t.data.(io.Closer); ok { - c.Close() - } - if sd, ok := td.(StatefulData); ok { - t.data = sd - } else { - t.data = &statelessDataWrapper{td, make([]bool, t.Info.NumPieces())} +func (t *torrent) setStorage(td Data) (err error) { + if t.data != nil { + t.data.Close() } + t.data = td return }