(tested) fix redis to pass tests

* delete info hash count key from redis (replaced with SCARD on infohash set)
* add GC test
* add peer.Addr() functio to always return unwrapped address if 4to6 appear
This commit is contained in:
Lawrence, Rendall
2022-04-15 01:29:57 +03:00
parent 5c2471ca9b
commit 397e106396
17 changed files with 287 additions and 261 deletions
+9 -12
View File
@@ -139,16 +139,18 @@ func New(provided Config) (storage.Storage, error) {
ps.wg.Add(1)
go func() {
defer ps.wg.Done()
t := time.NewTimer(cfg.GarbageCollectionInterval)
defer t.Stop()
for {
select {
case <-ps.closed:
return
case <-time.After(cfg.GarbageCollectionInterval):
case <-t.C:
before := time.Now().Add(-cfg.PeerLifetime)
log.Debug("storage: purging peers with no announces since", log.Fields{"before": before})
if err := ps.collectGarbage(before); err != nil {
log.Error(err)
}
start := time.Now()
ps.GC(before)
recordGCDuration(time.Since(start))
}
}
}()
@@ -546,20 +548,19 @@ func (ps *store) Delete(ctx string, keys ...string) error {
return nil
}
// collectGarbage deletes all Peers from the Storage which are older than the
// GC deletes all Peers from the Storage which are older than the
// cutoff time.
//
// This function must be able to execute while other methods on this interface
// are being executed in parallel.
func (ps *store) collectGarbage(cutoff time.Time) error {
func (ps *store) GC(cutoff time.Time) {
select {
case <-ps.closed:
return nil
return
default:
}
cutoffUnix := cutoff.UnixNano()
start := time.Now()
for _, shard := range ps.shards {
shard.RLock()
@@ -603,10 +604,6 @@ func (ps *store) collectGarbage(cutoff time.Time) error {
runtime.Gosched()
}
recordGCDuration(time.Since(start))
return nil
}
func (ps *store) Stop() stop.Result {
+50 -66
View File
@@ -2,24 +2,21 @@
// BitTorrent tracker keeping peer data in redis with hash.
// There two categories of hash:
//
// - CHI_{4,6}_{L,S}_infohash
// - CHI_{4,6}_{L,S}_<HASH> (hash type)
// To save peers that hold the infohash, used for fast searching,
// deleting, and timeout handling
//
// - CHI_{4,6}
// - CHI_{4,6}_I (set type)
// To save all the infohashes, used for garbage collection,
// metrics aggregation and leecher graduation
//
// Tree keys are used to record the count of swarms, seeders
// and leechers for each group (IPv4, IPv6).
//
// - CHI_{4,6}_I_C
// To record the number of infohashes.
//
// - CHI_{4,6}_S_C
// - CHI_{4,6}_S_C (key type)
// To record the number of seeders.
//
// - CHI_{4,6}_L_C
// - CHI_{4,6}_L_C (key type)
// To record the number of leechers.
package redis
@@ -64,8 +61,6 @@ const (
cnt6SeederKey = "CHI_6_C_S"
cnt4LeecherKey = "CHI_4_C_L"
cnt6LeecherKey = "CHI_6_C_L"
cnt4InfoHashKey = "CHI_4_C_I"
cnt6InfoHashKey = "CHI_6_C_I"
)
// ErrSentinelAndClusterChecked returned from initializer if both Config.Sentinel and Config.Cluster provided
@@ -277,10 +272,12 @@ func New(conf Config) (storage.Storage, error) {
ps.logFields = cfg.LogFields()
// Start a goroutine for garbage collection.
ps.wg.Add(1)
go ps.scheduleGC(cfg.GarbageCollectionInterval, cfg.PeerLifetime)
if cfg.PrometheusReportingInterval > 0 {
// Start a goroutine for reporting statistics to Prometheus.
ps.wg.Add(1)
go ps.schedulerProm(cfg.PrometheusReportingInterval)
} else {
log.Info("prometheus disabled because of zero reporting interval")
@@ -290,7 +287,6 @@ func New(conf Config) (storage.Storage, error) {
}
func (ps *store) scheduleGC(gcInterval, peerLifeTime time.Duration) {
ps.wg.Add(1)
defer ps.wg.Done()
t := time.NewTimer(gcInterval)
defer t.Stop()
@@ -299,12 +295,8 @@ func (ps *store) scheduleGC(gcInterval, peerLifeTime time.Duration) {
case <-ps.closed:
return
case <-t.C:
before := time.Now().Add(-peerLifeTime)
log.Debug("storage: purging peers with no announces since", log.Fields{"before": before})
cutoffUnix := before.UnixNano()
start := time.Now()
ps.gc(cutoffUnix, false)
ps.gc(cutoffUnix, true)
ps.GC(time.Now().Add(-peerLifeTime))
duration := time.Since(start).Milliseconds()
log.Debug("storage: recordGCDuration", log.Fields{"timeTaken(ms)": duration})
storage.PromGCDurationMilliseconds.Observe(float64(duration))
@@ -314,7 +306,6 @@ func (ps *store) scheduleGC(gcInterval, peerLifeTime time.Duration) {
}
func (ps *store) schedulerProm(reportInterval time.Duration) {
ps.wg.Add(1)
defer ps.wg.Done()
t := time.NewTicker(reportInterval)
for {
@@ -339,10 +330,16 @@ type store struct {
logFields log.Fields
}
func (ps *store) count(key string) (n uint64) {
func (ps *store) count(key string, getLength bool) (n uint64) {
var err error
if n, err = ps.con.Get(ps.ctx, key).Uint64(); err != nil && !errors.Is(err, redis.Nil) {
log.Error("storage: GET counter failure", log.Fields{
if getLength {
n, err = ps.con.SCard(ps.ctx, key).Uint64()
} else {
n, err = ps.con.Get(ps.ctx, key).Uint64()
}
err = asNil(err)
if err != nil {
log.Error("storage: len/counter failure", log.Fields{
"key": key,
"error": err,
})
@@ -355,15 +352,15 @@ func (ps *store) count(key string) (n uint64) {
func (ps *store) populateProm() {
numInfoHashes, numSeeders, numLeechers := new(uint64), new(uint64), new(uint64)
fetchFn := func(v6 bool) {
var cntSeederKey, cntLeecherKey, cntInfoHashKey string
var cntSeederKey, cntLeecherKey, ihSummaryKey string
if v6 {
cntSeederKey, cntLeecherKey, cntInfoHashKey = cnt6SeederKey, cnt6LeecherKey, cnt6InfoHashKey
cntSeederKey, cntLeecherKey, ihSummaryKey = cnt6SeederKey, cnt6LeecherKey, ih6Key
} else {
cntSeederKey, cntLeecherKey, cntInfoHashKey = cnt4SeederKey, cnt4LeecherKey, cnt4InfoHashKey
cntSeederKey, cntLeecherKey, ihSummaryKey = cnt4SeederKey, cnt4LeecherKey, ih4Key
}
*numInfoHashes += ps.count(cntInfoHashKey)
*numSeeders += ps.count(cntSeederKey)
*numLeechers += ps.count(cntLeecherKey)
*numInfoHashes += ps.count(ihSummaryKey, true)
*numSeeders += ps.count(cntSeederKey, false)
*numLeechers += ps.count(cntLeecherKey, false)
}
fetchFn(false)
@@ -403,15 +400,15 @@ func asNil(err error) error {
}
func (ps *store) PutSeeder(ih bittorrent.InfoHash, peer bittorrent.Peer) error {
var ihSummaryKey, ihPeerKey, cntPeerKey, cntInfoHashKey string
var ihSummaryKey, ihPeerKey, cntPeerKey string
log.Debug("storage: PutSeeder", log.Fields{
"InfoHash": ih,
"Peer": peer,
})
if peer.Addr().Is6() {
ihSummaryKey, ihPeerKey, cntPeerKey, cntInfoHashKey = ih6Key, ih6SeederKey, cnt6SeederKey, cnt6InfoHashKey
ihSummaryKey, ihPeerKey, cntPeerKey = ih6Key, ih6SeederKey, cnt6SeederKey
} else {
ihSummaryKey, ihPeerKey, cntPeerKey, cntInfoHashKey = ih4Key, ih4SeederKey, cnt4SeederKey, cnt4InfoHashKey
ihSummaryKey, ihPeerKey, cntPeerKey = ih4Key, ih4SeederKey, cnt4SeederKey
}
ihPeerKey += ih.RawString()
now := ps.getClock()
@@ -423,13 +420,7 @@ func (ps *store) PutSeeder(ih bittorrent.InfoHash, peer bittorrent.Peer) error {
if err = ps.con.Incr(ps.ctx, cntPeerKey).Err(); err != nil {
return
}
var added int64
if added, err = ps.con.SAdd(ps.ctx, ihSummaryKey, ihPeerKey).Result(); err != nil {
return
}
if added > 0 {
err = ps.con.Incr(ps.ctx, cntInfoHashKey).Err()
}
err = ps.con.SAdd(ps.ctx, ihSummaryKey, ihPeerKey).Err()
return
})
}
@@ -473,16 +464,14 @@ func (ps *store) PutLeecher(ih bittorrent.InfoHash, peer bittorrent.Peer) error
}
ihPeerKey += ih.RawString()
now := ps.getClock()
return ps.tx(func(tx redis.Pipeliner) (err error) {
if err = tx.HSet(ps.ctx, ihPeerKey, peer.RawString(), now).Err(); err != nil {
if err = tx.HSet(ps.ctx, ihPeerKey, peer.RawString(), ps.getClock()).Err(); err != nil {
return
}
if err = tx.Incr(ps.ctx, cntPeerKey).Err(); err != nil {
return err
}
err = tx.HSet(ps.ctx, ihSummaryKey, ihPeerKey, now).Err()
err = tx.SAdd(ps.ctx, ihSummaryKey, ihPeerKey).Err()
return
})
}
@@ -515,7 +504,7 @@ func (ps *store) DeleteLeecher(ih bittorrent.InfoHash, peer bittorrent.Peer) err
}
func (ps *store) GraduateLeecher(ih bittorrent.InfoHash, peer bittorrent.Peer) error {
var ihSummaryKey, ihSeederKey, ihLeecherKey, cntSeederKey, cntLeecherKey, cntInfoHashKey string
var ihSummaryKey, ihSeederKey, ihLeecherKey, cntSeederKey, cntLeecherKey string
log.Debug("storage: GraduateLeecher", log.Fields{
"InfoHash": ih,
"Peer": peer,
@@ -523,16 +512,14 @@ func (ps *store) GraduateLeecher(ih bittorrent.InfoHash, peer bittorrent.Peer) e
if peer.Addr().Is6() {
ihSummaryKey, ihSeederKey, cntSeederKey = ih6Key, ih6SeederKey, cnt6SeederKey
ihLeecherKey, cntLeecherKey, cntInfoHashKey = ih6LeecherKey, cnt6LeecherKey, cnt6InfoHashKey
ihLeecherKey, cntLeecherKey = ih6LeecherKey, cnt6LeecherKey
} else {
ihSummaryKey, ihSeederKey, cntSeederKey = ih4Key, ih4SeederKey, cnt4SeederKey
ihLeecherKey, cntLeecherKey, cntInfoHashKey = ih4LeecherKey, cnt4LeecherKey, cnt4InfoHashKey
ihLeecherKey, cntLeecherKey = ih4LeecherKey, cnt4LeecherKey
}
infoHash, peerKey := ih.RawString(), peer.RawString()
ihSeederKey, ihLeecherKey = ihSeederKey+infoHash, ihLeecherKey+infoHash
now := ps.getClock()
return ps.tx(func(tx redis.Pipeliner) error {
deleted, err := tx.HDel(ps.ctx, ihLeecherKey, peerKey).Uint64()
err = asNil(err)
@@ -542,16 +529,13 @@ func (ps *store) GraduateLeecher(ih bittorrent.InfoHash, peer bittorrent.Peer) e
}
}
if err == nil {
err = tx.HSet(ps.ctx, ihSeederKey, peerKey, now).Err()
err = tx.HSet(ps.ctx, ihSeederKey, peerKey, ps.getClock()).Err()
}
if err == nil {
err = tx.Incr(ps.ctx, cntSeederKey).Err()
}
if err == nil {
err = tx.HSet(ps.ctx, ihSummaryKey, ihSeederKey, now).Err()
}
if err == nil {
err = tx.Incr(ps.ctx, cntInfoHashKey).Err()
err = tx.SAdd(ps.ctx, ihSummaryKey, ihSeederKey).Err()
}
return err
})
@@ -567,7 +551,7 @@ func (ps *store) AnnouncePeers(ih bittorrent.InfoHash, seeder bool, numWant int,
})
if peer.Addr().Is6() {
ihSeederKey, ihLeecherKey = ih6SeederKey, cnt6LeecherKey
ihSeederKey, ihLeecherKey = ih6SeederKey, ih6LeecherKey
} else {
ihSeederKey, ihLeecherKey = ih4SeederKey, ih4LeecherKey
}
@@ -647,7 +631,7 @@ func (ps *store) ScrapeSwarm(ih bittorrent.InfoHash, peer bittorrent.Peer) (resp
})
resp.InfoHash = ih
if peer.Addr().Is6() {
ihSeederKey, ihLeecherKey = ih6SeederKey, cnt6LeecherKey
ihSeederKey, ihLeecherKey = ih6SeederKey, ih6LeecherKey
} else {
ihSeederKey, ihLeecherKey = ih4SeederKey, ih4LeecherKey
}
@@ -737,6 +721,13 @@ func (ps *store) Delete(ctx string, keys ...string) (err error) {
return
}
func (ps *store) GC(cutoff time.Time) {
log.Debug("storage: purging peers with no announces since", log.Fields{"before": cutoff})
cutoffUnix := cutoff.UnixNano()
ps.gc(cutoffUnix, false)
ps.gc(cutoffUnix, true)
}
// gc deletes all Peers from the Storage which are older than the
// cutoff time.
//
@@ -782,14 +773,14 @@ func (ps *store) Delete(ctx string, keys ...string) (err error) {
// - If the change happens after the HLEN, we will not even attempt to make the
// transaction. The infohash key will remain in the addressFamil hash and
// we'll attempt to clean it up the next time gc runs.
func (ps *store) gc(cutoffUnix int64, v6 bool) {
func (ps *store) gc(cutoffNanos int64, v6 bool) {
// list all infoHashKeys in the group
var ihSummaryKey, ihSeederKey, ihLeecherKey, cntSeederKey, cntLeecherKey, cntInfoHashKey string
var ihSummaryKey, ihSeederKey, ihLeecherKey, cntSeederKey, cntLeecherKey string
if v6 {
cntSeederKey, cntLeecherKey, cntInfoHashKey = cnt6SeederKey, cnt6LeecherKey, cnt6InfoHashKey
cntSeederKey, cntLeecherKey = cnt6SeederKey, cnt6LeecherKey
ihSummaryKey, ihSeederKey, ihLeecherKey = ih6Key, ih6SeederKey, ih6LeecherKey
} else {
cntSeederKey, cntLeecherKey, cntInfoHashKey = cnt4SeederKey, cnt4LeecherKey, cnt4InfoHashKey
cntSeederKey, cntLeecherKey = cnt4SeederKey, cnt4LeecherKey
ihSummaryKey, ihSeederKey, ihLeecherKey = ih4Key, ih4SeederKey, ih4LeecherKey
}
infoHashKeys, err := ps.con.SMembers(ps.ctx, ihSummaryKey).Result()
@@ -818,7 +809,7 @@ func (ps *store) gc(cutoffUnix int64, v6 bool) {
var peer bittorrent.Peer
if peer, err = bittorrent.NewPeer(peerKey); err == nil {
if mtime, err := strconv.ParseInt(timeStamp, 10, 64); err == nil {
if mtime <= cutoffUnix {
if mtime <= cutoffNanos {
log.Debug("storage: Redis: deleting peer", log.Fields{
"Peer": peer,
})
@@ -853,22 +844,15 @@ func (ps *store) gc(cutoffUnix int64, v6 bool) {
}
}
// use WATCH to avoid race condition
// https://redis.io/topics/transactions
err = asNil(ps.con.Watch(ps.ctx, func(tx *redis.Tx) (err error) {
var infoHashCount int64
infoHashCount, err = ps.con.HLen(ps.ctx, infoHashKey).Result()
var infoHashCount uint64
infoHashCount, err = ps.con.HLen(ps.ctx, infoHashKey).Uint64()
err = asNil(err)
if err == nil && infoHashCount == 0 {
// Empty hashes are not shown among existing keys,
// in other words, it's removed automatically after `HDEL` the last field.
// _, err := ps.con.Del(ps.ctx, infoHashKey)
var deletedCount int64
deletedCount, err = ps.con.SRem(ps.ctx, ihSummaryKey, infoHashKey).Result()
err = asNil(err)
if err == nil && seeder && deletedCount > 0 {
err = ps.con.Decr(ps.ctx, cntInfoHashKey).Err()
}
err = asNil(ps.con.SRem(ps.ctx, ihSummaryKey, infoHashKey).Err())
}
return err
}, infoHashKey))
+4
View File
@@ -3,6 +3,7 @@ package storage
import (
"errors"
"sync"
"time"
"github.com/sot-tech/mochi/bittorrent"
"github.com/sot-tech/mochi/pkg/log"
@@ -128,6 +129,9 @@ type Storage interface {
// Delete used to delete arbitrary data in specified context by its keys
Delete(ctx string, keys ...string) error
// GC used to delete stale data, such as timed out seeders/leechers
GC(cutoff time.Time)
// Stopper is an interface that expects a Stop method to stop the Storage.
// For more details see the documentation in the stop package.
stop.Stopper
+64 -68
View File
@@ -2,50 +2,46 @@ package test
import (
"math/rand"
"net"
"net/netip"
"runtime"
"sync/atomic"
"testing"
"github.com/sot-tech/mochi/bittorrent"
// used for seeding global math.Rand
_ "github.com/sot-tech/mochi/pkg/randseed"
"github.com/sot-tech/mochi/storage"
)
type benchData struct {
infohashes [1000]bittorrent.InfoHash
peers [1000]bittorrent.Peer
infoHashes [1000]bittorrent.InfoHash
peers [10000]bittorrent.Peer
}
func generateInfoHashes() (a [1000]bittorrent.InfoHash) {
for i := range a {
b := make([]byte, bittorrent.InfoHashV1Len)
rand.Read(b)
a[i], _ = bittorrent.NewInfoHash(b)
a[i] = randIH(rand.Int63()%2 == 0)
}
return
}
func generatePeers() (a [1000]bittorrent.Peer) {
func generatePeers() (a [10000]bittorrent.Peer) {
for i := range a {
ip := make([]byte, 4)
n, err := rand.Read(ip)
if err != nil || n != 4 {
panic("unable to create random bytes")
}
id := [bittorrent.PeerIDLen]byte{}
n, err = rand.Read(id[:])
if err != nil || n != bittorrent.InfoHashV1Len {
panic("unable to create random bytes")
var ip []byte
if rand.Int63()%2 == 0 {
ip = make([]byte, net.IPv4len)
} else {
ip = make([]byte, net.IPv6len)
}
rand.Read(ip)
addr, ok := netip.AddrFromSlice(ip)
if !ok {
panic("unable to create ip from random bytes")
}
port := uint16(rand.Uint32())
a[i] = bittorrent.Peer{
ID: id,
ID: randPeerID(),
AddrPort: netip.AddrPortFrom(addr, port),
}
}
@@ -122,7 +118,7 @@ func (bh *benchHolder) Nop(b *testing.B) {
// Put can run in parallel.
func (bh *benchHolder) Put(b *testing.B) {
bh.runBenchmark(b, true, nil, func(i int, ps storage.Storage, bd *benchData) error {
return ps.PutSeeder(bd.infohashes[0], bd.peers[0])
return ps.PutSeeder(bd.infoHashes[0], bd.peers[0])
})
}
@@ -132,27 +128,27 @@ func (bh *benchHolder) Put(b *testing.B) {
// Put1k can run in parallel.
func (bh *benchHolder) Put1k(b *testing.B) {
bh.runBenchmark(b, true, nil, func(i int, ps storage.Storage, bd *benchData) error {
return ps.PutSeeder(bd.infohashes[0], bd.peers[i%1000])
return ps.PutSeeder(bd.infoHashes[0], bd.peers[i%1000])
})
}
// Put1kInfoHash benchmarks the PutSeeder method of a storage.Storage by cycling
// through 1000 infohashes and putting the same peer into their swarms.
// through 1000 infoHashes and putting the same peer into their swarms.
//
// Put1kInfoHash can run in parallel.
func (bh *benchHolder) Put1kInfoHash(b *testing.B) {
bh.runBenchmark(b, true, nil, func(i int, ps storage.Storage, bd *benchData) error {
return ps.PutSeeder(bd.infohashes[i%1000], bd.peers[0])
return ps.PutSeeder(bd.infoHashes[i%1000], bd.peers[0])
})
}
// Put1kInfoHash1k benchmarks the PutSeeder method of a storage.Storage by cycling
// through 1000 infohashes and 1000 Peers and calling Put with them.
// through 1000 infoHashes and 1000 Peers and calling Put with them.
//
// Put1kInfoHash1k can run in parallel.
func (bh *benchHolder) Put1kInfoHash1k(b *testing.B) {
bh.runBenchmark(b, true, nil, func(i int, ps storage.Storage, bd *benchData) error {
err := ps.PutSeeder(bd.infohashes[i%1000], bd.peers[(i*3)%1000])
err := ps.PutSeeder(bd.infoHashes[i%1000], bd.peers[(i*3)%1000])
return err
})
}
@@ -163,11 +159,11 @@ func (bh *benchHolder) Put1kInfoHash1k(b *testing.B) {
// PutDelete can not run in parallel.
func (bh *benchHolder) PutDelete(b *testing.B) {
bh.runBenchmark(b, false, nil, func(i int, ps storage.Storage, bd *benchData) error {
err := ps.PutSeeder(bd.infohashes[0], bd.peers[0])
err := ps.PutSeeder(bd.infoHashes[0], bd.peers[0])
if err != nil {
return err
}
return ps.DeleteSeeder(bd.infohashes[0], bd.peers[0])
return ps.DeleteSeeder(bd.infoHashes[0], bd.peers[0])
})
}
@@ -177,39 +173,39 @@ func (bh *benchHolder) PutDelete(b *testing.B) {
// PutDelete1k can not run in parallel.
func (bh *benchHolder) PutDelete1k(b *testing.B) {
bh.runBenchmark(b, false, nil, func(i int, ps storage.Storage, bd *benchData) error {
err := ps.PutSeeder(bd.infohashes[0], bd.peers[i%1000])
err := ps.PutSeeder(bd.infoHashes[0], bd.peers[i%1000])
if err != nil {
return err
}
return ps.DeleteSeeder(bd.infohashes[0], bd.peers[i%1000])
return ps.DeleteSeeder(bd.infoHashes[0], bd.peers[i%1000])
})
}
// PutDelete1kInfoHash behaves like PutDelete1k with 1000 infohashes instead of
// PutDelete1kInfoHash behaves like PutDelete1k with 1000 infoHashes instead of
// 1000 Peers.
//
// PutDelete1kInfoHash can not run in parallel.
func (bh *benchHolder) PutDelete1kInfoHash(b *testing.B) {
bh.runBenchmark(b, false, nil, func(i int, ps storage.Storage, bd *benchData) error {
err := ps.PutSeeder(bd.infohashes[i%1000], bd.peers[0])
err := ps.PutSeeder(bd.infoHashes[i%1000], bd.peers[0])
if err != nil {
return err
}
return ps.DeleteSeeder(bd.infohashes[i%1000], bd.peers[0])
return ps.DeleteSeeder(bd.infoHashes[i%1000], bd.peers[0])
})
}
// PutDelete1kInfoHash1k behaves like PutDelete1k with 1000 infohashes in
// PutDelete1kInfoHash1k behaves like PutDelete1k with 1000 infoHashes in
// addition to 1000 Peers.
//
// PutDelete1kInfoHash1k can not run in parallel.
func (bh *benchHolder) PutDelete1kInfoHash1k(b *testing.B) {
bh.runBenchmark(b, false, nil, func(i int, ps storage.Storage, bd *benchData) error {
err := ps.PutSeeder(bd.infohashes[i%1000], bd.peers[(i*3)%1000])
err := ps.PutSeeder(bd.infoHashes[i%1000], bd.peers[(i*3)%1000])
if err != nil {
return err
}
err = ps.DeleteSeeder(bd.infohashes[i%1000], bd.peers[(i*3)%1000])
err = ps.DeleteSeeder(bd.infoHashes[i%1000], bd.peers[(i*3)%1000])
return err
})
}
@@ -220,7 +216,7 @@ func (bh *benchHolder) PutDelete1kInfoHash1k(b *testing.B) {
// DeleteNonexist can run in parallel.
func (bh *benchHolder) DeleteNonexist(b *testing.B) {
bh.runBenchmark(b, true, nil, func(i int, ps storage.Storage, bd *benchData) error {
_ = ps.DeleteSeeder(bd.infohashes[0], bd.peers[0])
_ = ps.DeleteSeeder(bd.infoHashes[0], bd.peers[0])
return nil
})
}
@@ -231,18 +227,18 @@ func (bh *benchHolder) DeleteNonexist(b *testing.B) {
// DeleteNonexist can run in parallel.
func (bh *benchHolder) DeleteNonexist1k(b *testing.B) {
bh.runBenchmark(b, true, nil, func(i int, ps storage.Storage, bd *benchData) error {
_ = ps.DeleteSeeder(bd.infohashes[0], bd.peers[i%1000])
_ = ps.DeleteSeeder(bd.infoHashes[0], bd.peers[i%1000])
return nil
})
}
// DeleteNonexist1kInfoHash benchmarks the DeleteSeeder method of a storage.Storage by
// attempting to delete one Peer from one of 1000 infohashes.
// attempting to delete one Peer from one of 1000 infoHashes.
//
// DeleteNonexist1kInfoHash can run in parallel.
func (bh *benchHolder) DeleteNonexist1kInfoHash(b *testing.B) {
bh.runBenchmark(b, true, nil, func(i int, ps storage.Storage, bd *benchData) error {
_ = ps.DeleteSeeder(bd.infohashes[i%1000], bd.peers[0])
_ = ps.DeleteSeeder(bd.infoHashes[i%1000], bd.peers[0])
return nil
})
}
@@ -253,7 +249,7 @@ func (bh *benchHolder) DeleteNonexist1kInfoHash(b *testing.B) {
// DeleteNonexist1kInfoHash1k can run in parallel.
func (bh *benchHolder) DeleteNonexist1kInfoHash1k(b *testing.B) {
bh.runBenchmark(b, true, nil, func(i int, ps storage.Storage, bd *benchData) error {
_ = ps.DeleteSeeder(bd.infohashes[i%1000], bd.peers[(i*3)%1000])
_ = ps.DeleteSeeder(bd.infoHashes[i%1000], bd.peers[(i*3)%1000])
return nil
})
}
@@ -264,7 +260,7 @@ func (bh *benchHolder) DeleteNonexist1kInfoHash1k(b *testing.B) {
// GradNonexist can run in parallel.
func (bh *benchHolder) GradNonexist(b *testing.B) {
bh.runBenchmark(b, true, nil, func(i int, ps storage.Storage, bd *benchData) error {
_ = ps.GraduateLeecher(bd.infohashes[0], bd.peers[0])
_ = ps.GraduateLeecher(bd.infoHashes[0], bd.peers[0])
return nil
})
}
@@ -275,7 +271,7 @@ func (bh *benchHolder) GradNonexist(b *testing.B) {
// GradNonexist1k can run in parallel.
func (bh *benchHolder) GradNonexist1k(b *testing.B) {
bh.runBenchmark(b, true, nil, func(i int, ps storage.Storage, bd *benchData) error {
_ = ps.GraduateLeecher(bd.infohashes[0], bd.peers[i%1000])
_ = ps.GraduateLeecher(bd.infoHashes[0], bd.peers[i%1000])
return nil
})
}
@@ -286,19 +282,19 @@ func (bh *benchHolder) GradNonexist1k(b *testing.B) {
// GradNonexist1kInfoHash can run in parallel.
func (bh *benchHolder) GradNonexist1kInfoHash(b *testing.B) {
bh.runBenchmark(b, true, nil, func(i int, ps storage.Storage, bd *benchData) error {
_ = ps.GraduateLeecher(bd.infohashes[i%1000], bd.peers[0])
_ = ps.GraduateLeecher(bd.infoHashes[i%1000], bd.peers[0])
return nil
})
}
// GradNonexist1kInfoHash1k benchmarks the GraduateLeecher method of a storage.Storage
// by attempting to graduate one of 1000 nonexistent Peers for one of 1000
// infohashes.
// infoHashes.
//
// GradNonexist1kInfoHash1k can run in parallel.
func (bh *benchHolder) GradNonexist1kInfoHash1k(b *testing.B) {
bh.runBenchmark(b, true, nil, func(i int, ps storage.Storage, bd *benchData) error {
_ = ps.GraduateLeecher(bd.infohashes[i%1000], bd.peers[(i*3)%1000])
_ = ps.GraduateLeecher(bd.infoHashes[i%1000], bd.peers[(i*3)%1000])
return nil
})
}
@@ -310,15 +306,15 @@ func (bh *benchHolder) GradNonexist1kInfoHash1k(b *testing.B) {
// PutGradDelete can not run in parallel.
func (bh *benchHolder) PutGradDelete(b *testing.B) {
bh.runBenchmark(b, false, nil, func(i int, ps storage.Storage, bd *benchData) error {
err := ps.PutLeecher(bd.infohashes[0], bd.peers[0])
err := ps.PutLeecher(bd.infoHashes[0], bd.peers[0])
if err != nil {
return err
}
err = ps.GraduateLeecher(bd.infohashes[0], bd.peers[0])
err = ps.GraduateLeecher(bd.infoHashes[0], bd.peers[0])
if err != nil {
return err
}
return ps.DeleteSeeder(bd.infohashes[0], bd.peers[0])
return ps.DeleteSeeder(bd.infoHashes[0], bd.peers[0])
})
}
@@ -327,51 +323,51 @@ func (bh *benchHolder) PutGradDelete(b *testing.B) {
// PutGradDelete1k can not run in parallel.
func (bh *benchHolder) PutGradDelete1k(b *testing.B) {
bh.runBenchmark(b, false, nil, func(i int, ps storage.Storage, bd *benchData) error {
err := ps.PutLeecher(bd.infohashes[0], bd.peers[i%1000])
err := ps.PutLeecher(bd.infoHashes[0], bd.peers[i%1000])
if err != nil {
return err
}
err = ps.GraduateLeecher(bd.infohashes[0], bd.peers[i%1000])
err = ps.GraduateLeecher(bd.infoHashes[0], bd.peers[i%1000])
if err != nil {
return err
}
return ps.DeleteSeeder(bd.infohashes[0], bd.peers[i%1000])
return ps.DeleteSeeder(bd.infoHashes[0], bd.peers[i%1000])
})
}
// PutGradDelete1kInfoHash behaves like PutGradDelete with one of 1000
// infohashes.
// infoHashes.
//
// PutGradDelete1kInfoHash can not run in parallel.
func (bh *benchHolder) PutGradDelete1kInfoHash(b *testing.B) {
bh.runBenchmark(b, false, nil, func(i int, ps storage.Storage, bd *benchData) error {
err := ps.PutLeecher(bd.infohashes[i%1000], bd.peers[0])
err := ps.PutLeecher(bd.infoHashes[i%1000], bd.peers[0])
if err != nil {
return err
}
err = ps.GraduateLeecher(bd.infohashes[i%1000], bd.peers[0])
err = ps.GraduateLeecher(bd.infoHashes[i%1000], bd.peers[0])
if err != nil {
return err
}
return ps.DeleteSeeder(bd.infohashes[i%1000], bd.peers[0])
return ps.DeleteSeeder(bd.infoHashes[i%1000], bd.peers[0])
})
}
// PutGradDelete1kInfoHash1k behaves like PutGradDelete with one of 1000 Peers
// and one of 1000 infohashes.
// and one of 1000 infoHashes.
//
// PutGradDelete1kInfoHash can not run in parallel.
func (bh *benchHolder) PutGradDelete1kInfoHash1k(b *testing.B) {
bh.runBenchmark(b, false, nil, func(i int, ps storage.Storage, bd *benchData) error {
err := ps.PutLeecher(bd.infohashes[i%1000], bd.peers[(i*3)%1000])
err := ps.PutLeecher(bd.infoHashes[i%1000], bd.peers[(i*3)%1000])
if err != nil {
return err
}
err = ps.GraduateLeecher(bd.infohashes[i%1000], bd.peers[(i*3)%1000])
err = ps.GraduateLeecher(bd.infoHashes[i%1000], bd.peers[(i*3)%1000])
if err != nil {
return err
}
err = ps.DeleteSeeder(bd.infohashes[i%1000], bd.peers[(i*3)%1000])
err = ps.DeleteSeeder(bd.infoHashes[i%1000], bd.peers[(i*3)%1000])
return err
})
}
@@ -381,9 +377,9 @@ func putPeers(ps storage.Storage, bd *benchData) error {
for j := 0; j < 1000; j++ {
var err error
if j < 1000/2 {
err = ps.PutLeecher(bd.infohashes[i], bd.peers[j])
err = ps.PutLeecher(bd.infoHashes[i], bd.peers[j])
} else {
err = ps.PutSeeder(bd.infohashes[i], bd.peers[j])
err = ps.PutSeeder(bd.infoHashes[i], bd.peers[j])
}
if err != nil {
return err
@@ -400,18 +396,18 @@ func putPeers(ps storage.Storage, bd *benchData) error {
// AnnounceLeecher can run in parallel.
func (bh *benchHolder) AnnounceLeecher(b *testing.B) {
bh.runBenchmark(b, true, putPeers, func(i int, ps storage.Storage, bd *benchData) error {
_, err := ps.AnnouncePeers(bd.infohashes[0], false, 50, bd.peers[0])
_, err := ps.AnnouncePeers(bd.infoHashes[0], false, 50, bd.peers[0])
return err
})
}
// AnnounceLeecher1kInfoHash behaves like AnnounceLeecher with one of 1000
// infohashes.
// infoHashes.
//
// AnnounceLeecher1kInfoHash can run in parallel.
func (bh *benchHolder) AnnounceLeecher1kInfoHash(b *testing.B) {
bh.runBenchmark(b, true, putPeers, func(i int, ps storage.Storage, bd *benchData) error {
_, err := ps.AnnouncePeers(bd.infohashes[i%1000], false, 50, bd.peers[0])
_, err := ps.AnnouncePeers(bd.infoHashes[i%1000], false, 50, bd.peers[0])
return err
})
}
@@ -422,18 +418,18 @@ func (bh *benchHolder) AnnounceLeecher1kInfoHash(b *testing.B) {
// AnnounceSeeder can run in parallel.
func (bh *benchHolder) AnnounceSeeder(b *testing.B) {
bh.runBenchmark(b, true, putPeers, func(i int, ps storage.Storage, bd *benchData) error {
_, err := ps.AnnouncePeers(bd.infohashes[0], true, 50, bd.peers[0])
_, err := ps.AnnouncePeers(bd.infoHashes[0], true, 50, bd.peers[0])
return err
})
}
// AnnounceSeeder1kInfoHash behaves like AnnounceSeeder with one of 1000
// infohashes.
// infoHashes.
//
// AnnounceSeeder1kInfoHash can run in parallel.
func (bh *benchHolder) AnnounceSeeder1kInfoHash(b *testing.B) {
bh.runBenchmark(b, true, putPeers, func(i int, ps storage.Storage, bd *benchData) error {
_, err := ps.AnnouncePeers(bd.infohashes[i%1000], true, 50, bd.peers[0])
_, err := ps.AnnouncePeers(bd.infoHashes[i%1000], true, 50, bd.peers[0])
return err
})
}
@@ -444,17 +440,17 @@ func (bh *benchHolder) AnnounceSeeder1kInfoHash(b *testing.B) {
// ScrapeSwarm can run in parallel.
func (bh *benchHolder) ScrapeSwarm(b *testing.B) {
bh.runBenchmark(b, true, putPeers, func(i int, ps storage.Storage, bd *benchData) error {
ps.ScrapeSwarm(bd.infohashes[0], bd.peers[0])
ps.ScrapeSwarm(bd.infoHashes[0], bd.peers[0])
return nil
})
}
// ScrapeSwarm1kInfoHash behaves like ScrapeSwarm with one of 1000 infohashes.
// ScrapeSwarm1kInfoHash behaves like ScrapeSwarm with one of 1000 infoHashes.
//
// ScrapeSwarm1kInfoHash can run in parallel.
func (bh *benchHolder) ScrapeSwarm1kInfoHash(b *testing.B) {
bh.runBenchmark(b, true, putPeers, func(i int, ps storage.Storage, bd *benchData) error {
ps.ScrapeSwarm(bd.infohashes[i%1000], bd.peers[0])
ps.ScrapeSwarm(bd.infoHashes[i%1000], bd.peers[0])
return nil
})
}
+16
View File
@@ -2,6 +2,7 @@ package test
import (
"testing"
"time"
"github.com/stretchr/testify/require"
@@ -244,6 +245,19 @@ func (th *testHolder) CustomBulkPutContainsLoadDelete(t *testing.T) {
}
}
func (th *testHolder) GC(t *testing.T) {
for _, c := range testData {
require.Nil(t, th.st.PutSeeder(c.ih, c.peer))
require.Nil(t, th.st.PutSeeder(c.ih, v4Peer))
require.Nil(t, th.st.PutSeeder(c.ih, v6Peer))
}
th.st.GC(time.Now().Add(time.Hour))
for _, c := range testData {
_, err := th.st.AnnouncePeers(c.ih, false, 100, v4Peer)
require.Equal(t, storage.ErrResourceDoesNotExist, err)
}
}
// RunTests tests a Storage implementation against the interface.
func RunTests(t *testing.T, p storage.Storage) {
th := testHolder{st: p}
@@ -275,6 +289,8 @@ func RunTests(t *testing.T, p storage.Storage) {
t.Run("CustomPutContainsLoadDelete", th.CustomPutContainsLoadDelete)
t.Run("CustomBulkPutContainsLoadDelete", th.CustomBulkPutContainsLoadDelete)
t.Run("GC", th.GC)
e := th.st.Stop()
require.Nil(t, <-e)
}
+28 -6
View File
@@ -1,9 +1,12 @@
package test
import (
"math/rand"
"net/netip"
"github.com/sot-tech/mochi/bittorrent"
// used for seeding global math.Rand
_ "github.com/sot-tech/mochi/pkg/randseed"
)
var (
@@ -13,13 +16,32 @@ var (
v4Peer, v6Peer bittorrent.Peer
)
func randIH(v2 bool) (ih bittorrent.InfoHash) {
var b []byte
if v2 {
b = make([]byte, bittorrent.InfoHashV2Len)
} else {
b = make([]byte, bittorrent.InfoHashV1Len)
}
rand.Read(b)
ih, _ = bittorrent.NewInfoHash(b)
return
}
func randPeerID() (ih bittorrent.PeerID) {
b := make([]byte, bittorrent.PeerIDLen)
rand.Read(b)
ih, _ = bittorrent.NewPeerID(b)
return
}
func init() {
testIh1, _ = bittorrent.NewInfoHash("00000000000000000001")
testIh2, _ = bittorrent.NewInfoHash("00000000000000000002")
testPeerID0, _ = bittorrent.NewPeerID([]byte("00000000000000000001"))
testPeerID1, _ = bittorrent.NewPeerID([]byte("00000000000000000002"))
testPeerID2, _ = bittorrent.NewPeerID([]byte("99999999999999999994"))
testPeerID3, _ = bittorrent.NewPeerID([]byte("99999999999999999996"))
testIh1 = randIH(false)
testIh2 = randIH(true)
testPeerID0 = randPeerID()
testPeerID1 = randPeerID()
testPeerID2 = randPeerID()
testPeerID3 = randPeerID()
testData = []hashPeer{
{
testIh1,