diff --git a/cmd/shisui/main.go b/cmd/shisui/main.go index 6977180a4dd9..35e754c09c00 100644 --- a/cmd/shisui/main.go +++ b/cmd/shisui/main.go @@ -127,6 +127,7 @@ func shisui(ctx *cli.Context) error { // Start system runtime metrics collection go metrics.CollectProcessMetrics(3 * time.Second) + go metrics.CollectPortalMetrics(5*time.Second, ctx.StringSlice(utils.PortalNetworksFlag.Name), ctx.String(utils.PortalDataDirFlag.Name)) if metrics.Enabled { storageCapacity = metrics.NewRegisteredGauge("portal/storage_capacity", nil) @@ -418,6 +419,7 @@ func initBeacon(config Config, server *rpc.Server, conn discover.UDPConn, localN DB: sqlDb, NodeId: localNode.ID(), Spec: configs.Mainnet, + NetworkName: portalwire.Beacon.Name(), }) if err != nil { return nil, err @@ -456,7 +458,7 @@ func initState(config Config, server *rpc.Server, conn discover.UDPConn, localNo if err != nil { return nil, err } - stateStore := state.NewStateStorage(contentStorage) + stateStore := state.NewStateStorage(contentStorage, db) contentQueue := make(chan *discover.ContentElement, 50) protocol, err := discover.NewPortalProtocol(config.Protocol, portalwire.State, config.PrivateKey, conn, localNode, discV5, stateStore, contentQueue) diff --git a/metrics/portal_metrics.go b/metrics/portal_metrics.go new file mode 100644 index 000000000000..8d524ffdddf3 --- /dev/null +++ b/metrics/portal_metrics.go @@ -0,0 +1,153 @@ +package metrics + +import ( + "database/sql" + "errors" + "os" + "path" + "slices" + "strings" + "time" + + "github.com/ethereum/go-ethereum/log" + "github.com/ethereum/go-ethereum/p2p/discover/portalwire" +) + +type networkFileMetric struct { + filename string + metric Gauge + file *os.File + network string +} + +type PortalStorageMetrics struct { + RadiusRatio GaugeFloat64 + EntriesCount Gauge + ContentStorageUsage Gauge +} + +const ( + countEntrySql = "SELECT COUNT(1) FROM kvstore;" + contentStorageUsageSql = "SELECT SUM( length(value) ) FROM kvstore;" +) + +// CollectPortalMetrics periodically collects various metrics about system entities. +func CollectPortalMetrics(refresh time.Duration, networks []string, dataDir string) { + // Short circuit if the metrics system is disabled + if !Enabled { + return + } + + // Define the various metrics to collect + var ( + historyTotalStorage = GetOrRegisterGauge("portal/history/total_storage", nil) + beaconTotalStorage = GetOrRegisterGauge("portal/beacon/total_storage", nil) + stateTotalStorage = GetOrRegisterGauge("portal/state/total_storage", nil) + ) + + var metricsArr []*networkFileMetric + if slices.Contains(networks, portalwire.History.Name()) { + dbPath := path.Join(dataDir, portalwire.History.Name()) + metricsArr = append(metricsArr, &networkFileMetric{ + filename: path.Join(dbPath, portalwire.History.Name()+".sqlite"), + metric: historyTotalStorage, + network: portalwire.History.Name(), + }) + } + if slices.Contains(networks, portalwire.Beacon.Name()) { + dbPath := path.Join(dataDir, portalwire.Beacon.Name()) + metricsArr = append(metricsArr, &networkFileMetric{ + filename: path.Join(dbPath, portalwire.Beacon.Name()+".sqlite"), + metric: beaconTotalStorage, + network: portalwire.Beacon.Name(), + }) + } + if slices.Contains(networks, portalwire.State.Name()) { + dbPath := path.Join(dataDir, portalwire.State.Name()) + metricsArr = append(metricsArr, &networkFileMetric{ + filename: path.Join(dbPath, portalwire.State.Name()+".sqlite"), + metric: stateTotalStorage, + network: portalwire.State.Name(), + }) + } + + for { + for _, m := range metricsArr { + var err error = nil + if m.file == nil { + m.file, err = os.OpenFile(m.filename, os.O_RDONLY, 0600) + if err != nil { + log.Debug("Could not open file", "network", m.network, "file", m.filename, "metric", "total_storage", "err", err) + } + } + if m.file != nil && err == nil { + stat, err := m.file.Stat() + if err != nil { + log.Warn("Could not get file stat", "network", m.network, "file", m.filename, "metric", "total_storage", "err", err) + } + if err == nil { + m.metric.Update(stat.Size()) + } + } + } + + time.Sleep(refresh) + } +} + +func NewPortalStorageMetrics(network string, db *sql.DB) (*PortalStorageMetrics, error) { + if !Enabled { + return nil, nil + } + + if network != portalwire.History.Name() && network != portalwire.Beacon.Name() && network != portalwire.State.Name() { + log.Debug("Unknow network for metrics", "network", network) + return nil, errors.New("unknow network for metrics") + } + + var countSql string + var contentSql string + if network == portalwire.Beacon.Name() { + countSql = strings.Replace(countEntrySql, "kvstore", "beacon", 1) + contentSql = strings.Replace(contentStorageUsageSql, "kvstore", "beacon", 1) + contentSql = strings.Replace(contentSql, "value", "content_value", 1) + } else { + countSql = countEntrySql + contentSql = contentStorageUsageSql + } + + storageMetrics := &PortalStorageMetrics{} + + storageMetrics.RadiusRatio = NewRegisteredGaugeFloat64("portal/"+network+"/radius_ratio", nil) + storageMetrics.RadiusRatio.Update(1) + + storageMetrics.EntriesCount = NewRegisteredGauge("portal/"+network+"/entry_count", nil) + log.Debug("Counting entities in " + network + " storage for metrics") + var res *int64 = new(int64) + q := db.QueryRow(countSql) + if q.Err() == sql.ErrNoRows { + storageMetrics.EntriesCount.Update(0) + } else if q.Err() != nil { + log.Error("Querry execution error", "network", network, "metric", "entry_count", "err", q.Err()) + return nil, q.Err() + } else { + q.Scan(res) + storageMetrics.EntriesCount.Update(*res) + } + + storageMetrics.ContentStorageUsage = NewRegisteredGauge("portal/"+network+"/content_storage", nil) + log.Debug("Counting storage usage (bytes) in " + network + " for metrics") + var res2 *int64 = new(int64) + q2 := db.QueryRow(contentSql) + if q2.Err() == sql.ErrNoRows { + storageMetrics.ContentStorageUsage.Update(0) + } else if q2.Err() != nil { + log.Error("Querry execution error", "network", network, "metric", "entry_count", "err", q2.Err()) + return nil, q2.Err() + } else { + q2.Scan(res2) + storageMetrics.ContentStorageUsage.Update(*res2) + } + + return storageMetrics, nil +} diff --git a/portalnetwork/beacon/storage.go b/portalnetwork/beacon/storage.go index a49a6135c012..ff90835313d8 100644 --- a/portalnetwork/beacon/storage.go +++ b/portalnetwork/beacon/storage.go @@ -6,6 +6,7 @@ import ( "database/sql" "github.com/ethereum/go-ethereum/log" + "github.com/ethereum/go-ethereum/metrics" "github.com/holiman/uint256" "github.com/protolambda/zrnt/eth2/beacon/common" "github.com/protolambda/ztyp/codec" @@ -23,6 +24,8 @@ type BeaconStorage struct { cache *beaconStorageCache } +var portalStorageMetrics *metrics.PortalStorageMetrics + type beaconStorageCache struct { OptimisticUpdate []byte FinalityUpdate []byte @@ -41,6 +44,13 @@ func NewBeaconStorage(config storage.PortalStorageConfig) (storage.ContentStorag if err := bs.setup(); err != nil { return nil, err } + + var err error + portalStorageMetrics, err = metrics.NewPortalStorageMetrics(config.NetworkName, config.DB) + if err != nil { + return nil, err + } + return bs, nil } @@ -175,10 +185,18 @@ func (bs *BeaconStorage) getLcUpdateValueByRange(start, end uint64) ([]byte, err func (bs *BeaconStorage) putContentValue(contentId, contentKey, value []byte) error { length := 32 + len(contentKey) + len(value) _, err := bs.db.ExecContext(context.Background(), InsertQueryBeacon, contentId, contentKey, value, length) + if metrics.Enabled && err != nil { + portalStorageMetrics.EntriesCount.Inc(1) + portalStorageMetrics.ContentStorageUsage.Inc(int64(len(value))) + } return err } func (bs *BeaconStorage) putLcUpdate(period uint64, value []byte) error { _, err := bs.db.ExecContext(context.Background(), InsertLCUpdateQuery, period, value, 0, len(value)) + if metrics.Enabled && err != nil { + portalStorageMetrics.EntriesCount.Inc(1) + portalStorageMetrics.ContentStorageUsage.Inc(int64(len(value))) + } return err } diff --git a/portalnetwork/history/storage.go b/portalnetwork/history/storage.go index 488d2ff3622c..ee43fb246eb5 100644 --- a/portalnetwork/history/storage.go +++ b/portalnetwork/history/storage.go @@ -21,10 +21,6 @@ import ( "github.com/mattn/go-sqlite3" ) -var ( - radiusRatio metrics.GaugeFloat64 -) - const ( sqliteName = "history.sqlite" contentDeletionFraction = 0.05 // 5% of the content will be deleted when the storage capacity is hit and radius gets adjusted. @@ -60,6 +56,8 @@ type ContentStorage struct { log log.Logger } +var portalStorageMetrics *metrics.PortalStorageMetrics + func xor(contentId, nodeId []byte) []byte { // length of contentId maybe not 32bytes padding := make([]byte, 32) @@ -112,10 +110,7 @@ func NewHistoryStorage(config storage.PortalStorageConfig) (storage.ContentStora log: log.New("storage", config.NetworkName), } hs.radius.Store(storage.MaxDistance) - if metrics.Enabled { - radiusRatio = metrics.NewRegisteredGaugeFloat64("portal/radius_ratio", nil) - radiusRatio.Update(1) - } + err := hs.createTable() if err != nil { return nil, err @@ -125,6 +120,14 @@ func NewHistoryStorage(config storage.PortalStorageConfig) (storage.ContentStora // Check whether we already have data, and use it to set radius + // necessary to test NetworkName==history because state also initialize HistoryStorage + if strings.ToLower(config.NetworkName) == "history" { + portalStorageMetrics, err = metrics.NewPortalStorageMetrics(config.NetworkName, config.DB) + if err != nil { + return nil, err + } + } + return hs, err } @@ -195,6 +198,10 @@ func (p *ContentStorage) put(contentId []byte, content []byte) PutResult { return PutResult{pruned: true, count: count} } + if metrics.Enabled { + portalStorageMetrics.EntriesCount.Inc(1) + portalStorageMetrics.ContentStorageUsage.Inc(int64(len(content))) + } return PutResult{} } @@ -289,6 +296,23 @@ func (p *ContentStorage) ContentSize() (uint64, error) { return p.queryRowUint64(sql) } +func (p *ContentStorage) SizeByKey(contentId []byte) (uint64, error) { + sql := "SELECT SUM( length(value) ) FROM kvstore WHERE key = " + string(contentId) + ";" + return p.queryRowUint64(sql) +} + +func (p *ContentStorage) SizeByKeys(ids [][]byte) (uint64, error) { + sql := "SELECT SUM( length(value) ) FROM kvstore WHERE key IN (?" + strings.Repeat(", ?", len(ids)-1) + ");" + return p.queryRowUint64(sql) +} + +func (p *ContentStorage) SizeOutRadius(radius *uint256.Int) (uint64, error) { + sql := "SELECT SUM( length(value) ) FROM kvstore WHERE greater(xor(key, (?1)), (?2)) = 1;" + var size uint64 + err := p.sqliteDB.QueryRow(sql, p.nodeId[:], radius.Bytes()).Scan(&size) + return size, err +} + func (p *ContentStorage) queryRowUint64(sqlStr string) (uint64, error) { // sql := "SELECT SUM(length(value)) FROM kvstore" stmt, err := p.sqliteDB.Prepare(sqlStr) @@ -345,7 +369,7 @@ func (p *ContentStorage) EstimateNewRadius(currentRadius *uint256.Int) (*uint256 newRadius := new(uint256.Int).Div(currentRadius, uint256.MustFromBig(bigFormat)) newRadius.Mul(newRadius, uint256.NewInt(100)) newRadius.Mod(newRadius, storage.MaxDistance) - radiusRatio.Update(newRadius.Float64() / 100) + portalStorageMetrics.RadiusRatio.Update(newRadius.Float64() / 100) } return new(uint256.Int).Div(currentRadius, uint256.MustFromBig(bigFormat)), nil } @@ -409,7 +433,7 @@ func (p *ContentStorage) deleteContentFraction(fraction float64) (deleteCount in if metrics.Enabled { dis.Mul(dis, uint256.NewInt(100)) dis.Mod(dis, storage.MaxDistance) - radiusRatio.Update(dis.Float64() / 100) + portalStorageMetrics.RadiusRatio.Update(dis.Float64() / 100) } } // row must close first, or database is locked @@ -423,11 +447,31 @@ func (p *ContentStorage) deleteContentFraction(fraction float64) (deleteCount in } func (p *ContentStorage) del(contentId []byte) error { - _, err := p.delStmt.Exec(contentId) + var sizeDel uint64 + var err error + if metrics.Enabled { + sizeDel, err = p.SizeByKey(contentId) + if err != nil { + return err + } + } + _, err = p.delStmt.Exec(contentId) + if metrics.Enabled && err != nil { + portalStorageMetrics.EntriesCount.Dec(1) + portalStorageMetrics.ContentStorageUsage.Dec(int64(sizeDel)) + } return err } func (p *ContentStorage) batchDel(ids [][]byte) error { + var sizeDel uint64 + var err error + if metrics.Enabled { + sizeDel, err = p.SizeByKeys(ids) + if err != nil { + return err + } + } query := "DELETE FROM kvstore WHERE key IN (?" + strings.Repeat(", ?", len(ids)-1) + ")" args := make([]interface{}, len(ids)) for i, id := range ids { @@ -435,7 +479,11 @@ func (p *ContentStorage) batchDel(ids [][]byte) error { } // delete items - _, err := p.sqliteDB.Exec(query, args...) + _, err = p.sqliteDB.Exec(query, args...) + if metrics.Enabled && err != nil { + portalStorageMetrics.EntriesCount.Dec(int64(len(args))) + portalStorageMetrics.ContentStorageUsage.Dec(int64(sizeDel)) + } return err } @@ -447,12 +495,24 @@ func (p *ContentStorage) ReclaimSpace() error { } func (p *ContentStorage) deleteContentOutOfRadius(radius *uint256.Int) error { + var sizeDel uint64 + var err error + if metrics.Enabled { + sizeDel, err = p.SizeOutRadius(radius) + if err != nil { + return err + } + } res, err := p.sqliteDB.Exec(deleteOutOfRadiusStmt, p.nodeId[:], radius.Bytes()) if err != nil { return err } - count, _ := res.RowsAffected() + count, err := res.RowsAffected() p.log.Trace("delete items", "count", count) + if metrics.Enabled && err != nil { + portalStorageMetrics.EntriesCount.Dec(count) + portalStorageMetrics.ContentStorageUsage.Dec(int64(sizeDel)) + } return err } diff --git a/portalnetwork/state/storage.go b/portalnetwork/state/storage.go index dae8a14d1da1..466c1707d6a2 100644 --- a/portalnetwork/state/storage.go +++ b/portalnetwork/state/storage.go @@ -3,10 +3,12 @@ package state import ( "bytes" "crypto/sha256" + "database/sql" "errors" "github.com/ethereum/go-ethereum/crypto" "github.com/ethereum/go-ethereum/log" + "github.com/ethereum/go-ethereum/metrics" "github.com/ethereum/go-ethereum/portalnetwork/storage" "github.com/holiman/uint256" "github.com/protolambda/ztyp/codec" @@ -21,14 +23,26 @@ var _ storage.ContentStorage = &StateStorage{} type StateStorage struct { store storage.ContentStorage + db *sql.DB log log.Logger } -func NewStateStorage(store storage.ContentStorage) *StateStorage { - return &StateStorage{ +var portalStorageMetrics *metrics.PortalStorageMetrics + +func NewStateStorage(store storage.ContentStorage, db *sql.DB) *StateStorage { + storage := &StateStorage{ store: store, + db: db, log: log.New("storage", "state"), } + + var err error + portalStorageMetrics, err = metrics.NewPortalStorageMetrics("state", db) + if err != nil { + return nil + } + + return storage } // Get implements storage.ContentStorage. @@ -84,6 +98,9 @@ func (s *StateStorage) putAccountTrieNode(contentKey []byte, contentId []byte, c err = s.store.Put(contentId, contentId, contentValueBuf.Bytes()) if err != nil { s.log.Error("failed to save data after validate", "type", contentKey[0], "key", contentKey[1:], "value", content) + } else if metrics.Enabled { + portalStorageMetrics.EntriesCount.Inc(1) + portalStorageMetrics.ContentStorageUsage.Inc(int64(len(content))) } return nil } @@ -118,6 +135,9 @@ func (s *StateStorage) putContractStorageTrieNode(contentKey []byte, contentId [ err = s.store.Put(contentId, contentId, contentValueBuf.Bytes()) if err != nil { s.log.Error("failed to save data after validate", "type", contentKey[0], "key", contentKey[1:], "value", content) + } else if metrics.Enabled { + portalStorageMetrics.EntriesCount.Inc(1) + portalStorageMetrics.ContentStorageUsage.Inc(int64(len(content))) } return nil } @@ -148,6 +168,9 @@ func (s *StateStorage) putContractBytecode(contentKey []byte, contentId []byte, err = s.store.Put(contentId, contentId, contentValueBuf.Bytes()) if err != nil { s.log.Error("failed to save data after validate", "type", contentKey[0], "key", contentKey[1:], "value", content) + } else if metrics.Enabled { + portalStorageMetrics.EntriesCount.Inc(1) + portalStorageMetrics.ContentStorageUsage.Inc(int64(len(content))) } return nil } diff --git a/portalnetwork/state/storage_test.go b/portalnetwork/state/storage_test.go index 7212e14407f7..180eaa9b0ada 100644 --- a/portalnetwork/state/storage_test.go +++ b/portalnetwork/state/storage_test.go @@ -10,7 +10,7 @@ import ( func TestStorage(t *testing.T) { storage := storage.NewMockStorage() - stateStorage := NewStateStorage(storage) + stateStorage := NewStateStorage(storage, nil) testfiles := []string{"account_trie_node.yaml", "contract_storage_trie_node.yaml", "contract_bytecode.yaml"} for _, file := range testfiles { cases, err := getTestCases(file)