Compare commits

...
13 Commits
Author SHA1 Message Date
Ian Norden efb334aef3 Merge pull request #51 from vulcanize/prerun
generate rand node id if none is provided
2021-10-28 10:52:35 -05:00
i-norden c83e8835ab generate rand node id if none is provided 2021-10-28 10:21:10 -05:00
S. R. Wadleigh 728642afec Merge pull request #48 from vulcanize/make_startup_script_executable
make startup_script.sh executable
2021-10-28 03:22:51 +00:00
Ian Norden 5875d1ec82 Update Dockerfile 2021-10-27 19:57:19 -05:00
erikdies 8d6919ab2b make startup_script.sh executable 2021-10-27 20:46:25 -04:00
Ian Norden abf72725ee Merge pull request #47 from vulcanize/prerun
cleanup
2021-10-27 17:46:13 -05:00
i-norden 8a55469403 cleanup 2021-10-27 17:45:47 -05:00
Ian Norden a3857bffe4 Merge pull request #46 from vulcanize/prerun
enable configuration of a prerun range by env variable; prerun only mode
2021-10-27 14:54:40 -05:00
i-norden 3956081874 sanitized prod config 2021-10-27 14:44:37 -05:00
i-norden e1abb399cd env vars for log stuff 2021-10-27 14:34:33 -05:00
i-norden 3d00e3ed05 enable configuration of a prerun range by env variable; prerun only mode 2021-10-27 13:36:19 -05:00
Ian Norden 84fff0aa03 Merge pull request #44 from vulcanize/stats
make DB stats collector registration optional
2021-10-25 18:57:40 -05:00
i-norden 79cb508d83 make DB stats collector registration optional 2021-10-25 18:54:54 -05:00
9 changed files with 188 additions and 18 deletions
+1 -1
View File
@@ -15,7 +15,7 @@ RUN GO111MODULE=on GCO_ENABLED=0 GOOS=linux go build -a -installsuffix cgo -ldfl
FROM alpine
ARG USER="vdm"
ARG CONFIG_FILE="./environments/example.toml"
ARG CONFIG_FILE="./environments/config.toml"
ARG EXPOSE_PORT=8545
RUN adduser -Du 5000 $USER adm
+23 -4
View File
@@ -29,10 +29,12 @@ const (
ETH_NETWORK_ID = "ETH_NETWORK_ID"
ETH_CHAIN_ID = "ETH_CHAIN_ID"
DB_CACHE_SIZE_MB = "DB_CACHE_SIZE_MB"
TRIE_CACHE_SIZE_MB = "TRIE_CACHE_SIZE_MB"
LVLDB_PATH = "LVLDB_PATH"
LVLDB_ANCIENT = "LVLDB_ANCIENT"
DB_CACHE_SIZE_MB = "DB_CACHE_SIZE_MB"
TRIE_CACHE_SIZE_MB = "TRIE_CACHE_SIZE_MB"
LVLDB_PATH = "LVLDB_PATH"
LVLDB_ANCIENT = "LVLDB_ANCIENT"
STATEDIFF_PRERUN = "STATEDIFF_PRERUN"
STATEDIFF_TRIE_WORKERS = "STATEDIFF_TRIE_WORKERS"
STATEDIFF_SERVICE_WORKERS = "STATEDIFF_SERVICE_WORKERS"
STATEDIFF_WORKER_QUEUE_SIZE = "STATEDIFF_WORKER_QUEUE_SIZE"
@@ -44,6 +46,14 @@ const (
PROM_HTTP = "PROM_HTTP"
PROM_HTTP_ADDR = "PROM_HTTP_ADDR"
PROM_HTTP_PORT = "PROM_HTTP_PORT"
PROM_DB_STATS = "PROM_DB_STATS"
PRERUN_ONLY = "PRERUN_ONLY"
PRERUN_RANGE_START = "PRERUN_RANGE_START"
PRERUN_RANGE_STOP = "PRERUN_RANGE_STOP"
LOG_LEVEL = "LOG_LEVEL"
LOG_FILE_PATH = "LOG_FILE_PATH"
)
// Bind env vars for eth node and DB configuration
@@ -76,8 +86,17 @@ func init() {
viper.BindEnv("prom.http", PROM_HTTP)
viper.BindEnv("prom.httpAddr", PROM_HTTP_ADDR)
viper.BindEnv("prom.httpPort", PROM_HTTP_PORT)
viper.BindEnv("prom.dbStats", PROM_DB_STATS)
viper.BindEnv("statediff.serviceWorkers", STATEDIFF_SERVICE_WORKERS)
viper.BindEnv("statediff.trieWorkers", STATEDIFF_TRIE_WORKERS)
viper.BindEnv("statediff.workerQueueSize", STATEDIFF_WORKER_QUEUE_SIZE)
viper.BindEnv("statediff.prerun", STATEDIFF_PRERUN)
viper.BindEnv("prerun.only", PRERUN_ONLY)
viper.BindEnv("prerun.start", PRERUN_RANGE_START)
viper.BindEnv("prerun.stop", PRERUN_RANGE_STOP)
viper.BindEnv("log.level", LOG_LEVEL)
viper.BindEnv("log.file", LOG_FILE_PATH)
}
+57 -9
View File
@@ -18,8 +18,10 @@ package cmd
import (
"fmt"
"math/rand"
"os"
"strings"
"time"
"github.com/ethereum/go-ethereum/statediff/indexer/node"
"github.com/ethereum/go-ethereum/statediff/indexer/postgres"
@@ -84,7 +86,6 @@ func initFuncs(cmd *cobra.Command, args []string) {
}
func logLevel() error {
viper.BindEnv("log.level", "LOGRUS_LEVEL")
lvl, err := log.ParseLevel(viper.GetString("log.level"))
if err != nil {
return err
@@ -136,8 +137,12 @@ func init() {
rootCmd.PersistentFlags().Bool("prom-http", false, "enable prometheus http service")
rootCmd.PersistentFlags().String("prom-http-addr", "127.0.0.1", "prometheus http host")
rootCmd.PersistentFlags().String("prom-http-port", "8080", "prometheus http port")
rootCmd.PersistentFlags().Bool("prom-db-stats", false, "enables prometheus db stats")
rootCmd.PersistentFlags().Bool("prom-metrics", false, "enable prometheus metrics")
rootCmd.PersistentFlags().Bool("metrics", false, "enable metrics")
rootCmd.PersistentFlags().Bool("prerun-only", false, "only process pre-configured ranges; exit afterwards")
rootCmd.PersistentFlags().Int("prerun-start", 0, "start height for a prerun range")
rootCmd.PersistentFlags().Int("prerun-stop", 0, "stop height for a prerun range")
viper.BindPFlag("server.httpPath", rootCmd.PersistentFlags().Lookup("http-path"))
viper.BindPFlag("server.ipcPath", rootCmd.PersistentFlags().Lookup("ipc-path"))
@@ -164,7 +169,13 @@ func init() {
viper.BindPFlag("prom.http", rootCmd.PersistentFlags().Lookup("prom-http"))
viper.BindPFlag("prom.httpAddr", rootCmd.PersistentFlags().Lookup("prom-http-addr"))
viper.BindPFlag("prom.httpPort", rootCmd.PersistentFlags().Lookup("prom-http-port"))
viper.BindPFlag("prom.metrics", rootCmd.PersistentFlags().Lookup("metrics"))
viper.BindPFlag("prom.dbStats", rootCmd.PersistentFlags().Lookup("prom-db-stats"))
viper.BindPFlag("prom.metrics", rootCmd.PersistentFlags().Lookup("prom-metrics"))
viper.BindPFlag("prerun.only", rootCmd.PersistentFlags().Lookup("prerun-only"))
viper.BindPFlag("prerun.start", rootCmd.PersistentFlags().Lookup("prerun-start"))
viper.BindPFlag("prerun.stop", rootCmd.PersistentFlags().Lookup("prerun-stop"))
rand.Seed(time.Now().UnixNano())
}
func initConfig() {
@@ -181,13 +192,50 @@ func initConfig() {
}
func GetEthNodeInfo() node.Info {
return node.Info{
ID: viper.GetString("ethereum.nodeID"),
ClientName: viper.GetString("ethereum.clientName"),
GenesisBlock: viper.GetString("ethereum.genesisBlock"),
NetworkID: viper.GetString("ethereum.networkID"),
ChainID: viper.GetUint64("ethereum.chainID"),
var nodeID, genesisBlock, networkID, clientName string
var chainID uint64
if !viper.IsSet("ethereum.nodeID") {
nodeID = randSeq(12)
} else {
nodeID = viper.GetString("ethereum.nodeID")
}
if !viper.IsSet("ethereum.genesisBlock") {
genesisBlock = "0xd4e56740f876aef8c010b86a40d5f56745a118d0906a34e69aec8c0db1cb8fa3"
} else {
genesisBlock = viper.GetString("ethereum.genesisBlock")
}
if !viper.IsSet("ethereum.chainID") {
chainID = 1
} else {
chainID = viper.GetUint64("ethereum.chainID")
}
if !viper.IsSet("ethereum.networkID") {
networkID = "1"
} else {
networkID = viper.GetString("ethereum.networkID")
}
if !viper.IsSet("ethereum.clientName") {
clientName = "eth-statediff-service"
} else {
clientName = viper.GetString("ethereum.clientName")
}
return node.Info{
ID: nodeID,
ClientName: clientName,
GenesisBlock: genesisBlock,
NetworkID: networkID,
ChainID: chainID,
}
}
var characters = []rune("abcdefghijklmnopqrstuvwxyzABCDEFGHIJKLMNOPQRSTUVWXYZ0123456789")
func randSeq(n int) string {
b := make([]rune, n)
for i := range b {
b[i] = characters[rand.Intn(len(characters))]
}
return string(b)
}
func GetDBParams() postgres.ConnectionParams {
+8
View File
@@ -55,6 +55,14 @@ func serve() {
logWithCommand.Fatal(err)
}
// short circuit if we only want to perform prerun
if viper.GetBool("prerun.only") {
if err := statediffService.Run(nil); err != nil {
logWithCommand.Fatal("unable to perform prerun: %v", err)
}
return
}
// start service and servers
logWithCommand.Info("Starting statediff service")
wg := new(sync.WaitGroup)
+20 -4
View File
@@ -52,15 +52,17 @@ func createStateDiffService() (sd.StateDiffService, error) {
}
// create statediff service
logWithCommand.Info("Creating statediff service")
logWithCommand.Info("Setting up Postgres DB")
db, err := setupPostgres(nodeInfo)
if err != nil {
logWithCommand.Fatal(err)
}
logWithCommand.Info("Creating statediff indexer")
indexer, err := ind.NewStateDiffIndexer(chainConf, db)
if err != nil {
logWithCommand.Fatal(err)
}
logWithCommand.Info("Creating statediff service")
sdConf := sd.Config{
ServiceWorkers: viper.GetUint("statediff.serviceWorkers"),
TrieWorkers: viper.GetUint("statediff.trieWorkers"),
@@ -71,12 +73,16 @@ func createStateDiffService() (sd.StateDiffService, error) {
}
func setupPostgres(nodeInfo node.Info) (*postgres.DB, error) {
params := GetDBParams()
db, err := postgres.NewDB(postgres.DbConnectionString(params), GetDBConfig(), nodeInfo)
p := GetDBParams()
logWithCommand.Info("initializing DB connection pool")
db, err := postgres.NewDB(postgres.DbConnectionString(p), GetDBConfig(), nodeInfo)
if err != nil {
return nil, err
}
prom.RegisterDBCollector(params.Name, db.DB)
if viper.GetBool("prom.dbStats") {
logWithCommand.Info("registering DB collector")
prom.RegisterDBCollector(p.Name, db.DB)
}
return db, nil
}
@@ -116,6 +122,16 @@ func setupPreRunRanges() []sd.RangeRequest {
Params: preRunParams,
}
}
if viper.IsSet("prerun.start") && viper.IsSet("prerun.stop") {
hardStart := viper.GetInt("prerun.start")
hardStop := viper.GetInt("prerun.stop")
blockRanges = append(blockRanges, sd.RangeRequest{
Start: uint64(hardStart),
Stop: uint64(hardStop),
Params: preRunParams,
})
}
return blockRanges
}
+51
View File
@@ -0,0 +1,51 @@
[leveldb]
path = "/app/geth-rw/chaindata"
ancient = "/app/geth-rw/chaindata/ancient"
[server]
ipcPath = ""
httpPath = "0.0.0.0:8545"
[statediff]
prerun = true
serviceWorkers = 1
workerQueueSize = 1024
trieWorkers = 16
[prerun]
only = true
ranges = []
[prerun.params]
intermediateStateNodes = true
intermediateStorageNodes = true
includeBlock = true
includeReceipts = true
includeTD = true
includeCode = true
watchedAddresses = []
watchedStorageKeys = []
[log]
file = ""
level = "info"
[eth]
chainID = 1
[database]
name = ""
hostname = ""
port = 5432
user = ""
password = ""
[cache]
database = 1024
trie = 4096
[prom]
dbStats = false
metrics = true
http = true
httpAddr = "0.0.0.0"
httpPort = 9100
+2
View File
@@ -13,6 +13,7 @@
trieWorkers = 4
[prerun]
only = false
ranges = [
[0, 1000]
]
@@ -45,6 +46,7 @@
trie = 1024
[prom]
dbStats = false
metrics = true
http = true
httpAddr = "localhost"
+26
View File
@@ -47,6 +47,8 @@ type StateDiffService interface {
Protocols() []p2p.Protocol
// Loop is the main event loop for processing state diffs
Loop(wg *sync.WaitGroup) error
// Run is a one-off command to run on a predefined set of ranges
Run(ranges []RangeRequest) error
// StateDiffAt method to get state diff object at specific block
StateDiffAt(blockNumber uint64, params sd.Params) (*sd.Payload, error)
// StateDiffFor method to get state diff object at specific block
@@ -115,11 +117,34 @@ func (sds *Service) APIs() []rpc.API {
}
}
// Run does a one-off processing run on the provided RangeRequests + any pre-runs, exiting afterwards
func (sds *Service) Run(rngs []RangeRequest) error {
for _, preRun := range sds.preruns {
logrus.Infof("processing prerun range (%d, %d)", preRun.Start, preRun.Stop)
for i := preRun.Start; i <= preRun.Stop; i++ {
if err := sds.WriteStateDiffAt(i, preRun.Params); err != nil {
return fmt.Errorf("error writing statediff at height %d in range (%d, %d) : %v", i, preRun.Start, preRun.Stop, err)
}
}
}
sds.preruns = nil
for _, rng := range rngs {
logrus.Infof("processing prerun range (%d, %d)", rng.Start, rng.Stop)
for i := rng.Start; i <= rng.Stop; i++ {
if err := sds.WriteStateDiffAt(i, rng.Params); err != nil {
return fmt.Errorf("error writing statediff at height %d in range (%d, %d) : %v", i, rng.Start, rng.Stop, err)
}
}
}
return nil
}
// Loop is an empty service loop for awaiting rpc requests
func (sds *Service) Loop(wg *sync.WaitGroup) error {
if sds.quitChan != nil {
return fmt.Errorf("service loop is already running")
}
sds.quitChan = make(chan struct{})
for i := 0; i < int(sds.workers); i++ {
wg.Add(1)
@@ -128,6 +153,7 @@ func (sds *Service) Loop(wg *sync.WaitGroup) error {
for {
select {
case blockRange := <-sds.queue:
logrus.Infof("service worker %d received range (%d, %d) off of work queue, beginning processing", id, blockRange.Start, blockRange.Stop)
prom.DecQueuedRanges()
for j := blockRange.Start; j <= blockRange.Stop; j++ {
if err := sds.WriteStateDiffAt(j, blockRange.Params); err != nil {
Regular → Executable
View File