DHT probe v0.0.1
This commit is contained in:
parent
84f7122a06
commit
60081acdfd
5 changed files with 459 additions and 175 deletions
|
|
@ -16,14 +16,17 @@ import (
|
|||
)
|
||||
|
||||
var (
|
||||
// errNoInputProvided indicates no input was passed
|
||||
errNoInputProvided = errors.New("no input provided")
|
||||
|
||||
// errInputIsNotAnURL indicates that input is not an URL
|
||||
errInputIsNotAnURL = errors.New("input is not an URL")
|
||||
|
||||
// errInvalidScheme indicates that the scheme is invalid
|
||||
errInvalidScheme = errors.New("scheme must be dht://")
|
||||
|
||||
// errInvalidPort indicates that no port was provided
|
||||
errInvalidPort = errors.New("no port was provided but dht:// requires explicit port")
|
||||
// errMissingPort indicates that no port was provided
|
||||
errMissingPort = errors.New("no port was provided but dht:// requires explicit port")
|
||||
)
|
||||
|
||||
const (
|
||||
|
|
@ -44,10 +47,14 @@ type RuntimeConfig struct {
|
|||
func config(input model.MeasurementTarget) (*RuntimeConfig, error) {
|
||||
// Bittorrent v2 hybrid test torrent: https://blog.libtorrent.org/2020/09/bittorrent-v2/
|
||||
// Has good chances of being seeded years from now
|
||||
hash := "631a31dd0a46257d5078c0dee4e66e26f73e42ac"
|
||||
hash := "631a31dd0a46257d5078c0dee4e66e26f73e42ac"
|
||||
|
||||
// TODO: static input from defaultDHTBoostrapNodes()
|
||||
// input == "" triggers runtime error from the experiment runner
|
||||
if input == "" {
|
||||
return nil, errNoInputProvided
|
||||
}
|
||||
|
||||
// TODO: static input from defaultDHTBoostrapNodes()
|
||||
// input == "" triggers runtime error from the experiment runner
|
||||
if input == "DUMMY" {
|
||||
// No requested DHT bootstrap node, let the DHT library try all it knows
|
||||
return &RuntimeConfig{
|
||||
|
|
@ -67,7 +74,7 @@ func config(input model.MeasurementTarget) (*RuntimeConfig, error) {
|
|||
|
||||
if parsed.Port() == "" {
|
||||
// Port is mandatory because DHT bootstrap nodes use different ports
|
||||
return nil, errInvalidPort
|
||||
return nil, errMissingPort
|
||||
}
|
||||
|
||||
valid_config := RuntimeConfig{
|
||||
|
|
@ -81,63 +88,63 @@ func config(input model.MeasurementTarget) (*RuntimeConfig, error) {
|
|||
|
||||
// TestKeys contains the experiment results
|
||||
type TestKeys struct {
|
||||
Queries []*model.ArchivalDNSLookupResult `json:"queries"`
|
||||
Runs []*IndividualTestKeys `json:"runs"`
|
||||
Queries []*model.ArchivalDNSLookupResult `json:"queries"`
|
||||
Runs []*IndividualTestKeys `json:"runs"`
|
||||
// Used for global failure (DNS resolution)
|
||||
Failure string `json:"failure"`
|
||||
// Indicated global or individual run failure
|
||||
Failed bool `json:"failed"`
|
||||
Failure string `json:"failure"`
|
||||
// Indicated global or individual run failure
|
||||
Failed bool `json:"failed"`
|
||||
}
|
||||
|
||||
func (tk *TestKeys) global_failure(err error) {
|
||||
tk.Failure = *tracex.NewFailure(err)
|
||||
tk.Failed = true
|
||||
tk.Failure = *tracex.NewFailure(err)
|
||||
tk.Failed = true
|
||||
}
|
||||
|
||||
func (tk *TestKeys) compute_global_failure() {
|
||||
if tk.Failure != "" {
|
||||
tk.Failed = true
|
||||
return
|
||||
}
|
||||
for _, itk := range(tk.Runs) {
|
||||
if itk.Failure != "" {
|
||||
tk.Failure = itk.Failure
|
||||
tk.Failed = true
|
||||
return
|
||||
}
|
||||
}
|
||||
if tk.Failure != "" {
|
||||
tk.Failed = true
|
||||
return
|
||||
}
|
||||
for _, itk := range tk.Runs {
|
||||
if itk.Failure != "" {
|
||||
tk.Failure = itk.Failure
|
||||
tk.Failed = true
|
||||
return
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Results for a single IP/port combo DHT bootstrap node
|
||||
// in case the DNS resolves to several IPs, or multiple bootstrap domains were used
|
||||
type IndividualTestKeys struct {
|
||||
// Logger, not exported to JSON
|
||||
logger model.Logger
|
||||
// Logger, not exported to JSON
|
||||
logger model.Logger
|
||||
|
||||
// List of IP/port combos tried to boostrap DHT
|
||||
BootstrapNodes []string `json:"bootstrap_nodes"`
|
||||
// Number of DHT bootsrap nodes
|
||||
BootstrapNum int `json:"bootstrap_num"`
|
||||
// Number of DHT peers contacted
|
||||
PeersTried uint32 `json:"peers_tried"`
|
||||
// Number of DHT peers who answered
|
||||
PeersResponded uint32 `json:"peers_responded"`
|
||||
// Number of DHT peers found for specific requested infohash
|
||||
InfohashPeers int `json:"infohash_peers"`
|
||||
// List of IP/port combos tried to boostrap DHT
|
||||
BootstrapNodes []string `json:"bootstrap_nodes"`
|
||||
// Number of DHT bootsrap nodes
|
||||
BootstrapNum int `json:"bootstrap_num"`
|
||||
// Number of DHT peers contacted
|
||||
PeersTried uint32 `json:"peers_tried"`
|
||||
// Number of DHT peers who answered
|
||||
PeersResponded uint32 `json:"peers_responded"`
|
||||
// Number of DHT peers found for specific requested infohash
|
||||
InfohashPeers int `json:"infohash_peers"`
|
||||
// Individual failure aborting the test run for this address/port combo
|
||||
Failure string `json:"failure"`
|
||||
Failure string `json:"failure"`
|
||||
}
|
||||
|
||||
func (itk *IndividualTestKeys) error(err error) {
|
||||
itk.Failure = *tracex.NewFailure(err)
|
||||
itk.logger.Warn(itk.Failure)
|
||||
itk.Failure = *tracex.NewFailure(err)
|
||||
itk.logger.Warn(itk.Failure)
|
||||
}
|
||||
|
||||
func NewITK(tk *TestKeys, log model.Logger) *IndividualTestKeys {
|
||||
itk := new(IndividualTestKeys)
|
||||
itk.logger = log
|
||||
tk.Runs = append(tk.Runs, itk)
|
||||
return itk
|
||||
itk := new(IndividualTestKeys)
|
||||
itk.logger = log
|
||||
tk.Runs = append(tk.Runs, itk)
|
||||
return itk
|
||||
}
|
||||
|
||||
type Measurer struct {
|
||||
|
|
@ -160,75 +167,75 @@ func (m Measurer) ExperimentVersion() string {
|
|||
}
|
||||
|
||||
func defaultDHTBoostrapNodes() []string {
|
||||
return []string{
|
||||
"router.utorrent.com:6881",
|
||||
"router.bittorrent.com:6881",
|
||||
"dht.transmissionbt.com:6881",
|
||||
"dht.aelitis.com:6881",
|
||||
"router.silotis.us:6881",
|
||||
"dht.libtorrent.org:25401",
|
||||
"dht.anacrolix.link:42069",
|
||||
"router.bittorrent.cloud:42069",
|
||||
}
|
||||
return []string{
|
||||
"router.utorrent.com:6881",
|
||||
"router.bittorrent.com:6881",
|
||||
"dht.transmissionbt.com:6881",
|
||||
"dht.aelitis.com:6881",
|
||||
"router.silotis.us:6881",
|
||||
"dht.libtorrent.org:25401",
|
||||
"dht.anacrolix.link:42069",
|
||||
"router.bittorrent.cloud:42069",
|
||||
}
|
||||
}
|
||||
|
||||
func DHTServer(bootstrap_nodes []string, itk *IndividualTestKeys) (*dht.Server, bool) {
|
||||
itk.BootstrapNodes = bootstrap_nodes
|
||||
itk.BootstrapNum = len(bootstrap_nodes)
|
||||
itk.BootstrapNodes = bootstrap_nodes
|
||||
itk.BootstrapNum = len(bootstrap_nodes)
|
||||
|
||||
// Starting new DHT client
|
||||
dhtconf := dht.NewDefaultServerConfig()
|
||||
dhtconf.QueryResendDelay = func() time.Duration {
|
||||
return 5 * time.Second
|
||||
}
|
||||
// Starting new DHT client
|
||||
dhtconf := dht.NewDefaultServerConfig()
|
||||
dhtconf.QueryResendDelay = func() time.Duration {
|
||||
return 5 * time.Second
|
||||
}
|
||||
|
||||
dhtconf.StartingNodes = func() (addrs []dht.Addr, err error) {
|
||||
for _, addrport := range(bootstrap_nodes) {
|
||||
udp_addr, err := net.ResolveUDPAddr("udp", addrport)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
addrs = append(addrs, dht.NewAddr(udp_addr))
|
||||
}
|
||||
return addrs, nil
|
||||
}
|
||||
dhtconf.StartingNodes = func() (addrs []dht.Addr, err error) {
|
||||
for _, addrport := range bootstrap_nodes {
|
||||
udp_addr, err := net.ResolveUDPAddr("udp", addrport)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
addrs = append(addrs, dht.NewAddr(udp_addr))
|
||||
}
|
||||
return addrs, nil
|
||||
}
|
||||
|
||||
dhtsrv, err := dht.NewServer(dhtconf)
|
||||
if err != nil {
|
||||
itk.error(err)
|
||||
return nil, false
|
||||
}
|
||||
itk.logger.Infof("Finished starting DHT server with bootstrap nodes: %v", bootstrap_nodes)
|
||||
return dhtsrv, true
|
||||
dhtsrv, err := dht.NewServer(dhtconf)
|
||||
if err != nil {
|
||||
itk.error(err)
|
||||
return nil, false
|
||||
}
|
||||
itk.logger.Infof("Finished starting DHT server with bootstrap nodes: %v", bootstrap_nodes)
|
||||
return dhtsrv, true
|
||||
}
|
||||
|
||||
func TestDHTServer(dht *dht.Server, infohash [20]byte, bootstrap_nodes []string, itk *IndividualTestKeys) bool {
|
||||
announce, err := dht.AnnounceTraversal(infohash)
|
||||
if err != nil {
|
||||
itk.error(err)
|
||||
return false
|
||||
}
|
||||
defer announce.Close()
|
||||
announce, err := dht.AnnounceTraversal(infohash)
|
||||
if err != nil {
|
||||
itk.error(err)
|
||||
return false
|
||||
}
|
||||
defer announce.Close()
|
||||
|
||||
counter := 0
|
||||
for entry := range announce.Peers {
|
||||
counter += 1
|
||||
itk.logger.Debugf("peer %d: %s", counter, entry.NodeInfo.Addr)
|
||||
}
|
||||
counter := 0
|
||||
for entry := range announce.Peers {
|
||||
counter += 1
|
||||
itk.logger.Debugf("peer %d: %s", counter, entry.NodeInfo.Addr)
|
||||
}
|
||||
|
||||
stats := announce.TraversalStats()
|
||||
itk.PeersTried = stats.NumAddrsTried
|
||||
itk.PeersResponded = stats.NumResponses
|
||||
itk.InfohashPeers = counter
|
||||
stats := announce.TraversalStats()
|
||||
itk.PeersTried = stats.NumAddrsTried
|
||||
itk.PeersResponded = stats.NumResponses
|
||||
itk.InfohashPeers = counter
|
||||
|
||||
if itk.PeersResponded == 0 {
|
||||
itk.error(errors.New("No DHT peers were found"))
|
||||
return false
|
||||
} else {
|
||||
itk.logger.Infof("Tried %d peers obtained from %d bootstrap nodes. Got response from %d. %d have requested infohash.", itk.PeersTried, itk.BootstrapNum, itk.PeersResponded, itk.InfohashPeers)
|
||||
}
|
||||
if itk.PeersResponded == 0 {
|
||||
itk.error(errors.New("No DHT peers were found"))
|
||||
return false
|
||||
} else {
|
||||
itk.logger.Infof("Tried %d peers obtained from %d bootstrap nodes. Got response from %d. %d have requested infohash.", itk.PeersTried, itk.BootstrapNum, itk.PeersResponded, itk.InfohashPeers)
|
||||
}
|
||||
|
||||
return true
|
||||
return true
|
||||
|
||||
}
|
||||
|
||||
|
|
@ -252,68 +259,68 @@ func (m Measurer) Run(
|
|||
ctx, cancel := context.WithTimeout(ctx, 60*time.Second)
|
||||
defer cancel()
|
||||
|
||||
// Turn string infohash into 20-bytes array
|
||||
var infohash [20]byte
|
||||
copy(infohash[:], config.infohash)
|
||||
// Turn string infohash into 20-bytes array
|
||||
var infohash [20]byte
|
||||
copy(infohash[:], config.infohash)
|
||||
|
||||
resolver := trace.NewStdlibResolver(log)
|
||||
resolver := trace.NewStdlibResolver(log)
|
||||
|
||||
if config.dhtnode != "" {
|
||||
// Specific node provided: resolve it
|
||||
log.Infof("Resolving DNS for %s", config.dhtnode)
|
||||
resolved_addrs, err := resolver.LookupHost(ctx, config.dhtnode)
|
||||
tk.Queries = append(tk.Queries, trace.DNSLookupsFromRoundTrip()...)
|
||||
if err != nil {
|
||||
tk.global_failure(err)
|
||||
return nil
|
||||
}
|
||||
log.Infof("Finished DNS for %s: %v", config.dhtnode, resolved_addrs)
|
||||
if config.dhtnode != "" {
|
||||
// Specific node provided: resolve it
|
||||
log.Infof("Resolving DNS for %s", config.dhtnode)
|
||||
resolved_addrs, err := resolver.LookupHost(ctx, config.dhtnode)
|
||||
tk.Queries = append(tk.Queries, trace.DNSLookupsFromRoundTrip()...)
|
||||
if err != nil {
|
||||
tk.global_failure(err)
|
||||
return nil
|
||||
}
|
||||
log.Infof("Finished DNS for %s: %v", config.dhtnode, resolved_addrs)
|
||||
|
||||
for _, addr := range(resolved_addrs) {
|
||||
for _, addr := range resolved_addrs {
|
||||
|
||||
node_addrport := net.JoinHostPort(addr, config.port)
|
||||
log.Infof("Trying DHT bootstrap node %s", node_addrport)
|
||||
node_addrports := []string{node_addrport}
|
||||
node_addrport := net.JoinHostPort(addr, config.port)
|
||||
log.Infof("Trying DHT bootstrap node %s", node_addrport)
|
||||
node_addrports := []string{node_addrport}
|
||||
|
||||
itk := NewITK(tk, log)
|
||||
itk := NewITK(tk, log)
|
||||
|
||||
dht, success := DHTServer(node_addrports, itk)
|
||||
if ! success {
|
||||
continue
|
||||
}
|
||||
dht, success := DHTServer(node_addrports, itk)
|
||||
if !success {
|
||||
continue
|
||||
}
|
||||
|
||||
TestDHTServer(dht, infohash, node_addrports, itk)
|
||||
}
|
||||
} else {
|
||||
// Use default DHT bootstrap nodes because none was given by input
|
||||
resolved_addrports := []string{}
|
||||
for _, bootstrap_domain := range(defaultDHTBoostrapNodes()) {
|
||||
// Ignore error because we use static input so panic chance is 0
|
||||
host, port, _ := net.SplitHostPort(bootstrap_domain)
|
||||
log.Infof("Resolving DNS for %s", host)
|
||||
resolved_addrs, err := resolver.LookupHost(ctx, host)
|
||||
tk.Queries = append(tk.Queries, trace.DNSLookupsFromRoundTrip()...)
|
||||
if err != nil {
|
||||
tk.global_failure(err)
|
||||
return nil
|
||||
}
|
||||
log.Infof("Finished DNS for %s: %v", host, resolved_addrs)
|
||||
for _, resolved_addr := range(resolved_addrs) {
|
||||
resolved_addrports = append(resolved_addrports, net.JoinHostPort(resolved_addr, port))
|
||||
}
|
||||
}
|
||||
log.Infof("Resolved the following bootstrap nodes: %v", resolved_addrports)
|
||||
TestDHTServer(dht, infohash, node_addrports, itk)
|
||||
}
|
||||
} else {
|
||||
// Use default DHT bootstrap nodes because none was given by input
|
||||
resolved_addrports := []string{}
|
||||
for _, bootstrap_domain := range defaultDHTBoostrapNodes() {
|
||||
// Ignore error because we use static input so panic chance is 0
|
||||
host, port, _ := net.SplitHostPort(bootstrap_domain)
|
||||
log.Infof("Resolving DNS for %s", host)
|
||||
resolved_addrs, err := resolver.LookupHost(ctx, host)
|
||||
tk.Queries = append(tk.Queries, trace.DNSLookupsFromRoundTrip()...)
|
||||
if err != nil {
|
||||
tk.global_failure(err)
|
||||
return nil
|
||||
}
|
||||
log.Infof("Finished DNS for %s: %v", host, resolved_addrs)
|
||||
for _, resolved_addr := range resolved_addrs {
|
||||
resolved_addrports = append(resolved_addrports, net.JoinHostPort(resolved_addr, port))
|
||||
}
|
||||
}
|
||||
log.Infof("Resolved the following bootstrap nodes: %v", resolved_addrports)
|
||||
|
||||
itk := NewITK(tk, log)
|
||||
dht, success := DHTServer(resolved_addrports, itk)
|
||||
if success {
|
||||
TestDHTServer(dht, infohash, resolved_addrports, itk)
|
||||
}
|
||||
}
|
||||
itk := NewITK(tk, log)
|
||||
dht, success := DHTServer(resolved_addrports, itk)
|
||||
if success {
|
||||
TestDHTServer(dht, infohash, resolved_addrports, itk)
|
||||
}
|
||||
}
|
||||
|
||||
tk.compute_global_failure()
|
||||
tk.compute_global_failure()
|
||||
|
||||
return nil
|
||||
return nil
|
||||
}
|
||||
|
||||
// NewExperimentMeasurer creates a new ExperimentMeasurer.
|
||||
|
|
|
|||
125
internal/engine/experiment/dht/dht_test.go
Normal file
125
internal/engine/experiment/dht/dht_test.go
Normal file
|
|
@ -0,0 +1,125 @@
|
|||
package dht
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"fmt"
|
||||
"log"
|
||||
"testing"
|
||||
|
||||
"github.com/anacrolix/dht/v2"
|
||||
"github.com/ooni/probe-cli/v3/internal/engine/mockable"
|
||||
"github.com/ooni/probe-cli/v3/internal/model"
|
||||
)
|
||||
|
||||
func TestMeasurer_run(t *testing.T) {
|
||||
// expectedPings is the expected number of pings
|
||||
const expectedPings = 4
|
||||
|
||||
// runHelper is an helper function to run this set of tests.
|
||||
runHelper := func(input string) (*model.Measurement, model.ExperimentMeasurer, error) {
|
||||
m := NewExperimentMeasurer(Config{})
|
||||
if m.ExperimentName() != "dht" {
|
||||
t.Fatal("invalid experiment name")
|
||||
}
|
||||
if m.ExperimentVersion() != "0.0.1" {
|
||||
t.Fatal("invalid experiment version")
|
||||
}
|
||||
ctx := context.Background()
|
||||
meas := &model.Measurement{
|
||||
Input: model.MeasurementTarget(input),
|
||||
}
|
||||
sess := &mockable.Session{
|
||||
MockableLogger: model.DiscardLogger,
|
||||
}
|
||||
callbacks := model.NewPrinterCallbacks(model.DiscardLogger)
|
||||
err := m.Run(ctx, sess, meas, callbacks)
|
||||
return meas, m, err
|
||||
}
|
||||
|
||||
t.Run("with empty input", func(t *testing.T) {
|
||||
_, _, err := runHelper("")
|
||||
if !errors.Is(err, errNoInputProvided) {
|
||||
t.Fatal("unexpected error", err)
|
||||
}
|
||||
})
|
||||
|
||||
t.Run("with invalid URL", func(t *testing.T) {
|
||||
_, _, err := runHelper("\t")
|
||||
if !errors.Is(err, errInputIsNotAnURL) {
|
||||
t.Fatal("unexpected error", err)
|
||||
}
|
||||
})
|
||||
|
||||
t.Run("with invalid scheme", func(t *testing.T) {
|
||||
_, _, err := runHelper("https://8.8.8.8:443/")
|
||||
if !errors.Is(err, errInvalidScheme) {
|
||||
t.Fatal("unexpected error", err)
|
||||
}
|
||||
})
|
||||
|
||||
t.Run("with missing port", func(t *testing.T) {
|
||||
_, _, err := runHelper("dht://8.8.8.8")
|
||||
if !errors.Is(err, errMissingPort) {
|
||||
t.Fatal("unexpected error", err)
|
||||
}
|
||||
})
|
||||
|
||||
t.Run("with local listener", func(t *testing.T) {
|
||||
conf := new(dht.ServerConfig)
|
||||
conf.StartingNodes = func() (addrs []dht.Addr, err error) {
|
||||
return []dht.Addr{}, nil
|
||||
}
|
||||
conf.Passive = false
|
||||
dht, err := dht.NewServer(conf)
|
||||
if err != nil {
|
||||
log.Fatal(err)
|
||||
}
|
||||
defer dht.Close()
|
||||
_, _ = dht.Bootstrap()
|
||||
|
||||
println(dht.Addr().String())
|
||||
url := fmt.Sprintf("dht://%s", dht.Addr().String())
|
||||
|
||||
meas, m, err := runHelper(url)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
tk := meas.TestKeys.(*TestKeys)
|
||||
|
||||
if tk.Failure != "" {
|
||||
t.Fatal(tk.Failure)
|
||||
}
|
||||
|
||||
if len(tk.Runs) != 1 {
|
||||
t.Fatal("Expected one DHT run")
|
||||
}
|
||||
|
||||
run := tk.Runs[0]
|
||||
if run.Failure != "" {
|
||||
t.Fatal(run.Failure)
|
||||
}
|
||||
|
||||
if run.BootstrapNum != 1 {
|
||||
t.Fatal("Expected only one bootstrap node")
|
||||
}
|
||||
|
||||
if run.PeersResponded != 1 {
|
||||
t.Fatal("Expected bootstrap node to respond")
|
||||
}
|
||||
|
||||
if tk.Failed == true {
|
||||
t.Fatal("Found TestKeys.Failed but found no global/run failure")
|
||||
}
|
||||
|
||||
ask, err := m.GetSummaryKeys(meas)
|
||||
if err != nil {
|
||||
t.Fatal("cannot obtain summary")
|
||||
}
|
||||
summary := ask.(SummaryKeys)
|
||||
if summary.IsAnomaly {
|
||||
t.Fatal("expected no anomaly")
|
||||
}
|
||||
})
|
||||
}
|
||||
22
internal/registry/dht.go
Normal file
22
internal/registry/dht.go
Normal file
|
|
@ -0,0 +1,22 @@
|
|||
package registry
|
||||
|
||||
//
|
||||
// Registers the `dnsping' experiment.
|
||||
//
|
||||
|
||||
import (
|
||||
"github.com/ooni/probe-cli/v3/internal/engine/experiment/dht"
|
||||
"github.com/ooni/probe-cli/v3/internal/model"
|
||||
)
|
||||
|
||||
func init() {
|
||||
AllExperiments["dht"] = &Factory{
|
||||
build: func(config interface{}) model.ExperimentMeasurer {
|
||||
return dht.NewExperimentMeasurer(
|
||||
*config.(*dht.Config),
|
||||
)
|
||||
},
|
||||
config: &dht.Config{},
|
||||
inputPolicy: model.InputOrStaticDefault,
|
||||
}
|
||||
}
|
||||
Loading…
Reference in a new issue