diff --git a/cmd/devp2p/crawl.go b/cmd/devp2p/crawl.go
new file mode 100644
index 0000000000..ae14425503
--- /dev/null
+++ b/cmd/devp2p/crawl.go
@@ -0,0 +1,151 @@
+// Copyright 2019 The go-ethereum Authors
+// This file is part of go-ethereum.
+//
+// go-ethereum is free software: you can redistribute it and/or modify
+// it under the terms of the GNU General Public License as published by
+// the Free Software Foundation, either version 3 of the License, or
+// (at your option) any later version.
+//
+// go-ethereum is distributed in the hope that it will be useful,
+// but WITHOUT ANY WARRANTY; without even the implied warranty of
+// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
+// GNU General Public License for more details.
+//
+// You should have received a copy of the GNU General Public License
+// along with go-ethereum. If not, see .
+
+package main
+
+import (
+ "time"
+
+ "github.com/ethereum/go-ethereum/log"
+ "github.com/ethereum/go-ethereum/p2p/discover"
+ "github.com/ethereum/go-ethereum/p2p/enode"
+)
+
+type crawler struct {
+ input nodeSet
+ output nodeSet
+ disc *discover.UDPv4
+ iters []enode.Iterator
+ inputIter enode.Iterator
+ ch chan *enode.Node
+ closed chan struct{}
+
+ // settings
+ revalidateInterval time.Duration
+}
+
+func newCrawler(input nodeSet, disc *discover.UDPv4, iters ...enode.Iterator) *crawler {
+ c := &crawler{
+ input: input,
+ output: make(nodeSet, len(input)),
+ disc: disc,
+ iters: iters,
+ inputIter: enode.IterNodes(input.nodes()),
+ ch: make(chan *enode.Node),
+ closed: make(chan struct{}),
+ }
+ c.iters = append(c.iters, c.inputIter)
+ // Copy input to output initially. Any nodes that fail validation
+ // will be dropped from output during the run.
+ for id, n := range input {
+ c.output[id] = n
+ }
+ return c
+}
+
+func (c *crawler) run(timeout time.Duration) nodeSet {
+ var (
+ timeoutTimer = time.NewTimer(timeout)
+ timeoutCh <-chan time.Time
+ doneCh = make(chan enode.Iterator, len(c.iters))
+ liveIters = len(c.iters)
+ )
+ for _, it := range c.iters {
+ go c.runIterator(doneCh, it)
+ }
+
+loop:
+ for {
+ select {
+ case n := <-c.ch:
+ c.updateNode(n)
+ case it := <-doneCh:
+ liveIters--
+ if liveIters == 0 {
+ break loop
+ }
+ if it == c.inputIter {
+ // Enable timeout when we're done revalidating the input nodes.
+ log.Info("Revalidation of input set is done", "len", len(c.input))
+ if timeout > 0 {
+ timeoutCh = timeoutTimer.C
+ }
+ }
+ case <-timeoutCh:
+ break loop
+ }
+ }
+
+ close(c.closed)
+ for _, it := range c.iters {
+ it.Close()
+ }
+ for ; liveIters > 0; liveIters-- {
+ <-doneCh
+ }
+ return c.output
+}
+
+func (c *crawler) runIterator(done chan<- enode.Iterator, it enode.Iterator) {
+ defer func() { done <- it }()
+ for it.Next() {
+ select {
+ case c.ch <- it.Node():
+ case <-c.closed:
+ return
+ }
+ }
+}
+
+func (c *crawler) updateNode(n *enode.Node) {
+ existing, ok := c.output[n.ID()]
+
+ // Skip validation of recently-seen nodes.
+ if ok && time.Since(existing.LastSeen) < c.revalidateInterval {
+ return
+ }
+
+ // Request the node record.
+ nn, err := c.disc.RequestENR(n)
+ if err != nil {
+ if existing.Checks == 0 {
+ log.Debug("Skipping node", "id", n.ID())
+ return
+ }
+ existing.Checks /= 2
+ } else {
+ if !ok {
+ existing.FirstSeen = truncNow()
+ }
+ existing.N = nn
+ existing.Seq = nn.Seq()
+ existing.LastSeen = truncNow()
+ existing.Checks++
+ }
+
+ // Store/update node in output set.
+ if existing.Checks <= 0 {
+ log.Info("Removing node", "id", n.ID())
+ delete(c.output, n.ID())
+ } else {
+ log.Info("Updating node", "id", n.ID(), "seq", existing.Seq, "checks", existing.Checks)
+ c.output[n.ID()] = existing
+ }
+}
+
+func truncNow() time.Time {
+ return time.Now().UTC().Truncate(1 * time.Second)
+}
diff --git a/cmd/devp2p/discv4cmd.go b/cmd/devp2p/discv4cmd.go
index ab5b874029..9525bec668 100644
--- a/cmd/devp2p/discv4cmd.go
+++ b/cmd/devp2p/discv4cmd.go
@@ -39,6 +39,7 @@ var (
discv4RequestRecordCommand,
discv4ResolveCommand,
discv4ResolveJSONCommand,
+ discv4CrawlCommand,
},
}
discv4PingCommand = cli.Command{
@@ -67,12 +68,25 @@ var (
Flags: []cli.Flag{bootnodesFlag},
ArgsUsage: "",
}
+ discv4CrawlCommand = cli.Command{
+ Name: "crawl",
+ Usage: "Updates a nodes.json file with random nodes found in the DHT",
+ Action: discv4Crawl,
+ Flags: []cli.Flag{bootnodesFlag, crawlTimeoutFlag},
+ }
)
-var bootnodesFlag = cli.StringFlag{
- Name: "bootnodes",
- Usage: "Comma separated nodes used for bootstrapping",
-}
+var (
+ bootnodesFlag = cli.StringFlag{
+ Name: "bootnodes",
+ Usage: "Comma separated nodes used for bootstrapping",
+ }
+ crawlTimeoutFlag = cli.DurationFlag{
+ Name: "timeout",
+ Usage: "Time limit for the crawl.",
+ Value: 30 * time.Minute,
+ }
+)
func discv4Ping(ctx *cli.Context) error {
n := getNodeArg(ctx)
@@ -113,30 +127,48 @@ func discv4ResolveJSON(ctx *cli.Context) error {
if ctx.NArg() < 1 {
return fmt.Errorf("need nodes file as argument")
}
- disc := startV4(ctx)
- defer disc.Close()
- file := ctx.Args().Get(0)
-
- // Load existing nodes in file.
- var nodes []*enode.Node
- if common.FileExist(file) {
- nodes = loadNodesJSON(file).nodes()
+ nodesFile := ctx.Args().Get(0)
+ inputSet := make(nodeSet)
+ if common.FileExist(nodesFile) {
+ inputSet = loadNodesJSON(nodesFile)
}
- // Add nodes from command line arguments.
+
+ // Add extra nodes from command line arguments.
+ var nodeargs []*enode.Node
for i := 1; i < ctx.NArg(); i++ {
n, err := parseNode(ctx.Args().Get(i))
if err != nil {
exit(err)
}
- nodes = append(nodes, n)
+ nodeargs = append(nodeargs, n)
}
- result := make(nodeSet, len(nodes))
- for _, n := range nodes {
- n = disc.Resolve(n)
- result[n.ID()] = nodeJSON{Seq: n.Seq(), N: n}
+ // Run the crawler.
+ disc := startV4(ctx)
+ defer disc.Close()
+ c := newCrawler(inputSet, disc, enode.IterNodes(nodeargs))
+ c.revalidateInterval = 0
+ output := c.run(0)
+ writeNodesJSON(nodesFile, output)
+ return nil
+}
+
+func discv4Crawl(ctx *cli.Context) error {
+ if ctx.NArg() < 1 {
+ return fmt.Errorf("need nodes file as argument")
}
- writeNodesJSON(file, result)
+ nodesFile := ctx.Args().First()
+ var inputSet nodeSet
+ if common.FileExist(nodesFile) {
+ inputSet = loadNodesJSON(nodesFile)
+ }
+
+ disc := startV4(ctx)
+ defer disc.Close()
+ c := newCrawler(inputSet, disc, disc.RandomNodes())
+ c.revalidateInterval = 10 * time.Minute
+ output := c.run(ctx.Duration(crawlTimeoutFlag.Name))
+ writeNodesJSON(nodesFile, output)
return nil
}
diff --git a/cmd/devp2p/dnscmd.go b/cmd/devp2p/dnscmd.go
index 74d70d3aaa..b510b4e3dc 100644
--- a/cmd/devp2p/dnscmd.go
+++ b/cmd/devp2p/dnscmd.go
@@ -109,7 +109,8 @@ func dnsSync(ctx *cli.Context) error {
}
def := treeToDefinition(url, t)
def.Meta.LastModified = time.Now()
- writeTreeDefinition(outdir, def)
+ writeTreeMetadata(outdir, def)
+ writeTreeNodes(outdir, def)
return nil
}
@@ -151,7 +152,7 @@ func dnsSign(ctx *cli.Context) error {
def = treeToDefinition(url, t)
def.Meta.LastModified = time.Now()
- writeTreeDefinition(defdir, def)
+ writeTreeMetadata(defdir, def)
return nil
}
@@ -315,26 +316,28 @@ func ensureValidTreeSignature(t *dnsdisc.Tree, pubkey *ecdsa.PublicKey, sig stri
return nil
}
-// writeTreeDefinition writes a DNS node tree definition to the given directory.
-func writeTreeDefinition(directory string, def *dnsDefinition) {
+// writeTreeMetadata writes a DNS node tree metata file to the given directory.
+func writeTreeMetadata(directory string, def *dnsDefinition) {
metaJSON, err := json.MarshalIndent(&def.Meta, "", jsonIndent)
if err != nil {
exit(err)
}
- // Convert nodes.
- nodes := make(nodeSet, len(def.Nodes))
- nodes.add(def.Nodes...)
- // Write.
if err := os.Mkdir(directory, 0744); err != nil && !os.IsExist(err) {
exit(err)
}
- metaFile, nodesFile := treeDefinitionFiles(directory)
- writeNodesJSON(nodesFile, nodes)
+ metaFile, _ := treeDefinitionFiles(directory)
if err := ioutil.WriteFile(metaFile, metaJSON, 0644); err != nil {
exit(err)
}
}
+func writeTreeNodes(directory string, def *dnsDefinition) {
+ ns := make(nodeSet, len(def.Nodes))
+ ns.add(def.Nodes...)
+ _, nodesFile := treeDefinitionFiles(directory)
+ writeNodesJSON(nodesFile, ns)
+}
+
func treeDefinitionFiles(directory string) (string, string) {
meta := filepath.Join(directory, "enrtree-info.json")
nodes := filepath.Join(directory, "nodes.json")
diff --git a/cmd/devp2p/main.go b/cmd/devp2p/main.go
index c88fe6f612..6faa650937 100644
--- a/cmd/devp2p/main.go
+++ b/cmd/devp2p/main.go
@@ -60,6 +60,7 @@ func init() {
enrdumpCommand,
discv4Command,
dnsCommand,
+ nodesetCommand,
}
}
diff --git a/cmd/devp2p/nodeset.go b/cmd/devp2p/nodeset.go
index a4a05016e9..80678d635c 100644
--- a/cmd/devp2p/nodeset.go
+++ b/cmd/devp2p/nodeset.go
@@ -22,6 +22,7 @@ import (
"fmt"
"io/ioutil"
"sort"
+ "time"
"github.com/ethereum/go-ethereum/common"
"github.com/ethereum/go-ethereum/p2p/enode"
@@ -34,8 +35,11 @@ const jsonIndent = " "
type nodeSet map[enode.ID]nodeJSON
type nodeJSON struct {
- Seq uint64 `json:"seq"`
- N *enode.Node `json:"record"`
+ Seq uint64 `json:"seq"`
+ N *enode.Node `json:"record"`
+ FirstSeen time.Time `json:"firstSeen,omitempty"`
+ LastSeen time.Time `json:"lastSeen,omitempty"`
+ Checks int `json:"checks"`
}
func loadNodesJSON(file string) nodeSet {
@@ -70,7 +74,7 @@ func (ns nodeSet) nodes() []*enode.Node {
func (ns nodeSet) add(nodes ...*enode.Node) {
for _, n := range nodes {
- ns[n.ID()] = nodeJSON{Seq: n.Seq(), N: n}
+ ns[n.ID()] = nodeJSON{Seq: n.Seq(), N: n, FirstSeen: truncNow()}
}
}
diff --git a/cmd/devp2p/nodesetcmd.go b/cmd/devp2p/nodesetcmd.go
new file mode 100644
index 0000000000..57db301d92
--- /dev/null
+++ b/cmd/devp2p/nodesetcmd.go
@@ -0,0 +1,49 @@
+// Copyright 2019 The go-ethereum Authors
+// This file is part of go-ethereum.
+//
+// go-ethereum is free software: you can redistribute it and/or modify
+// it under the terms of the GNU General Public License as published by
+// the Free Software Foundation, either version 3 of the License, or
+// (at your option) any later version.
+//
+// go-ethereum is distributed in the hope that it will be useful,
+// but WITHOUT ANY WARRANTY; without even the implied warranty of
+// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
+// GNU General Public License for more details.
+//
+// You should have received a copy of the GNU General Public License
+// along with go-ethereum. If not, see .
+
+package main
+
+import (
+ "fmt"
+
+ "gopkg.in/urfave/cli.v1"
+)
+
+var (
+ nodesetCommand = cli.Command{
+ Name: "nodeset",
+ Usage: "Node set tools",
+ Subcommands: []cli.Command{
+ nodesetInfoCommand,
+ },
+ }
+ nodesetInfoCommand = cli.Command{
+ Name: "info",
+ Usage: "Shows statistics about a node set",
+ Action: nodesetInfo,
+ ArgsUsage: "",
+ }
+)
+
+func nodesetInfo(ctx *cli.Context) error {
+ if ctx.NArg() < 1 {
+ return fmt.Errorf("need nodes file as argument")
+ }
+
+ ns := loadNodesJSON(ctx.Args().First())
+ fmt.Printf("Set contains %d nodes.\n", len(ns))
+ return nil
+}