mirror of
https://gitlab.com/pulsechaincom/prysm-pulse.git
synced 2025-01-15 22:48:19 +00:00
ade61717a4
* first version * cli context * fix service * starting change to ccache * ristretto cache * added test * test on evict * remove evict test * test onevict * comment for exported flag * update all span maps on load * fix setup db * span cache added to help flags * start save cache on exit * save cache to db before close * comment fix * fix flags * setup db new * data update from archive node * gaz * slashing detection on old attestations * un-export * rename * nishant feedback * workspace cr * lint fix * fix calls * start db * fix test * Update slasher/db/db.go Co-Authored-By: Nishant Das <nishdas93@gmail.com> * add flag * fix fail to start beacon client * mock beacon service * fix imports * gaz * goimports * add clear db flag * print finalized epoch * better msg * Update slasher/db/attester_slashings.go Co-Authored-By: Raul Jordan <raul@prysmaticlabs.com> * raul feedback * raul feedback * raul feedback * raul feedback * raul feedback * add detection in runtime * fix tests * raul feedbacks * raul feedback * raul feedback * goimports * Update beacon-chain/blockchain/process_attestation_helpers.go * Update beacon-chain/blockchain/receive_block.go * Update beacon-chain/core/blocks/block_operations_test.go * Update beacon-chain/core/blocks/block_operations.go * Update beacon-chain/core/epoch/epoch_processing.go * Update beacon-chain/sync/validate_aggregate_proof_test.go * Update shared/testutil/block.go * Update slasher/service/data_update.go * Update tools/blocktree/main.go * Update slasher/service/service.go * Update beacon-chain/core/epoch/precompute/attestation_test.go * Update beacon-chain/core/helpers/committee_test.go * Update beacon-chain/core/state/transition_test.go * Update beacon-chain/rpc/aggregator/server_test.go * Update beacon-chain/sync/validate_aggregate_proof.go * Update beacon-chain/rpc/validator/proposer_test.go * Update beacon-chain/blockchain/forkchoice/process_attestation.go Co-Authored-By: terence tsao <terence@prysmaticlabs.com> * Update slasher/db/indexed_attestations.go Co-Authored-By: terence tsao <terence@prysmaticlabs.com> * Update slasher/service/data_update.go Co-Authored-By: terence tsao <terence@prysmaticlabs.com> * terence feedback * terence feedback * goimports Co-authored-by: prylabs-bulldozer[bot] <58059840+prylabs-bulldozer[bot]@users.noreply.github.com> Co-authored-by: Raul Jordan <raul@prysmaticlabs.com> Co-authored-by: Nishant Das <nish1993@hotmail.com> Co-authored-by: terence tsao <terence@prysmaticlabs.com>
315 lines
10 KiB
Go
315 lines
10 KiB
Go
package db
|
|
|
|
import (
|
|
"bytes"
|
|
"fmt"
|
|
"reflect"
|
|
"sort"
|
|
|
|
"github.com/boltdb/bolt"
|
|
"github.com/gogo/protobuf/proto"
|
|
"github.com/pkg/errors"
|
|
ethpb "github.com/prysmaticlabs/ethereumapis/eth/v1alpha1"
|
|
slashpb "github.com/prysmaticlabs/prysm/proto/slashing"
|
|
"github.com/prysmaticlabs/prysm/shared/bytesutil"
|
|
"github.com/prysmaticlabs/prysm/shared/hashutil"
|
|
"github.com/prysmaticlabs/prysm/shared/params"
|
|
)
|
|
|
|
func createIndexedAttestation(enc []byte) (*ethpb.IndexedAttestation, error) {
|
|
protoIdxAtt := ðpb.IndexedAttestation{}
|
|
err := proto.Unmarshal(enc, protoIdxAtt)
|
|
if err != nil {
|
|
return nil, errors.Wrap(err, "failed to unmarshal encoding")
|
|
}
|
|
return protoIdxAtt, nil
|
|
}
|
|
|
|
func createValidatorIDsToIndexedAttestationList(enc []byte) (*slashpb.ValidatorIDToIdxAttList, error) {
|
|
protoIdxAtt := &slashpb.ValidatorIDToIdxAttList{}
|
|
err := proto.Unmarshal(enc, protoIdxAtt)
|
|
if err != nil {
|
|
return nil, errors.Wrap(err, "failed to unmarshal encoding")
|
|
}
|
|
return protoIdxAtt, nil
|
|
}
|
|
|
|
// IndexedAttestation accepts a epoch and validator index and returns a list of
|
|
// indexed attestations.
|
|
// Returns nil if the indexed attestation does not exist.
|
|
func (db *Store) IndexedAttestation(targetEpoch uint64, validatorID uint64) ([]*ethpb.IndexedAttestation, error) {
|
|
var iAtt []*ethpb.IndexedAttestation
|
|
key := bytesutil.Bytes8(targetEpoch)
|
|
err := db.view(func(tx *bolt.Tx) error {
|
|
bucket := tx.Bucket(indexedAttestationsIndicesBucket)
|
|
enc := bucket.Get(key)
|
|
iList, err := createValidatorIDsToIndexedAttestationList(enc)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
for _, a := range iList.IndicesList {
|
|
i := sort.Search(len(a.Indices), func(i int) bool { return a.Indices[i] >= validatorID })
|
|
if i < len(a.Indices) && a.Indices[i] == validatorID {
|
|
iaBucket := tx.Bucket(historicIndexedAttestationsBucket)
|
|
key := encodeEpochSig(targetEpoch, a.Signature)
|
|
enc = iaBucket.Get(key)
|
|
if len(enc) == 0 {
|
|
continue
|
|
}
|
|
iA, err := createIndexedAttestation(enc)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
iAtt = append(iAtt, iA)
|
|
}
|
|
}
|
|
return nil
|
|
})
|
|
|
|
return iAtt, err
|
|
}
|
|
|
|
// IndexedAttestations accepts a target epoch and returns a list of
|
|
// indexed attestations.
|
|
// Returns nil if the indexed attestation does not exist with that target epoch.
|
|
func (db *Store) IndexedAttestations(targetEpoch uint64) ([]*ethpb.IndexedAttestation, error) {
|
|
var iAtt []*ethpb.IndexedAttestation
|
|
key := bytesutil.Bytes8(targetEpoch)
|
|
err := db.view(func(tx *bolt.Tx) error {
|
|
c := tx.Bucket(historicIndexedAttestationsBucket).Cursor()
|
|
for k, enc := c.Seek(key); k != nil && bytes.Equal(k[:8], key); k, _ = c.Next() {
|
|
iA, err := createIndexedAttestation(enc)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
iAtt = append(iAtt, iA)
|
|
}
|
|
return nil
|
|
})
|
|
return iAtt, err
|
|
}
|
|
|
|
// LatestIndexedAttestationsTargetEpoch returns latest target epoch in db
|
|
// returns 0 if there is no indexed attestations in db.
|
|
func (db *Store) LatestIndexedAttestationsTargetEpoch() (uint64, error) {
|
|
var lt uint64
|
|
err := db.view(func(tx *bolt.Tx) error {
|
|
c := tx.Bucket(historicIndexedAttestationsBucket).Cursor()
|
|
k, _ := c.Last()
|
|
if k == nil {
|
|
return nil
|
|
}
|
|
lt = bytesutil.FromBytes8(k[:8])
|
|
return nil
|
|
})
|
|
return lt, err
|
|
}
|
|
|
|
// LatestValidatorIdx returns latest validator id in db
|
|
// returns 0 if there is no validators in db.
|
|
func (db *Store) LatestValidatorIdx() (uint64, error) {
|
|
var lt uint64
|
|
err := db.view(func(tx *bolt.Tx) error {
|
|
c := tx.Bucket(indexedAttestationsIndicesBucket).Cursor()
|
|
k, _ := c.Last()
|
|
if k == nil {
|
|
return nil
|
|
}
|
|
lt = bytesutil.FromBytes8(k[:8])
|
|
return nil
|
|
})
|
|
return lt, err
|
|
}
|
|
|
|
// DoubleVotes looks up db for slashable attesting data that were preformed by the same validator.
|
|
func (db *Store) DoubleVotes(targetEpoch uint64, validatorIdx uint64, dataRoot []byte, origAtt *ethpb.IndexedAttestation) ([]*ethpb.AttesterSlashing, error) {
|
|
idxAttestations, err := db.IndexedAttestation(targetEpoch, validatorIdx)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
if idxAttestations == nil || len(idxAttestations) == 0 {
|
|
return nil, fmt.Errorf("can't check nil indexed attestation for double vote")
|
|
}
|
|
var slashIdxAtt []*ethpb.IndexedAttestation
|
|
for _, at := range idxAttestations {
|
|
root, err := hashutil.HashProto(at.Data)
|
|
if at.Data == nil {
|
|
continue
|
|
}
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
if !bytes.Equal(root[:], dataRoot) {
|
|
slashIdxAtt = append(slashIdxAtt, at)
|
|
}
|
|
}
|
|
var as []*ethpb.AttesterSlashing
|
|
for _, ia := range slashIdxAtt {
|
|
as = append(as, ðpb.AttesterSlashing{
|
|
Attestation_1: origAtt,
|
|
Attestation_2: ia,
|
|
})
|
|
}
|
|
return as, nil
|
|
}
|
|
|
|
// HasIndexedAttestation accepts an epoch and validator id and returns true if the indexed attestation exists.
|
|
func (db *Store) HasIndexedAttestation(targetEpoch uint64, validatorID uint64) bool {
|
|
key := bytesutil.Bytes8(targetEpoch)
|
|
var hasAttestation bool
|
|
// #nosec G104
|
|
_ = db.view(func(tx *bolt.Tx) error {
|
|
bucket := tx.Bucket(indexedAttestationsIndicesBucket)
|
|
enc := bucket.Get(key)
|
|
iList, err := createValidatorIDsToIndexedAttestationList(enc)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
for _, a := range iList.IndicesList {
|
|
i := sort.Search(len(a.Indices), func(i int) bool { return a.Indices[i] >= validatorID })
|
|
if i < len(a.Indices) && a.Indices[i] == validatorID {
|
|
hasAttestation = true
|
|
return nil
|
|
}
|
|
}
|
|
hasAttestation = false
|
|
return nil
|
|
})
|
|
|
|
return hasAttestation
|
|
}
|
|
|
|
// SaveIndexedAttestation accepts epoch and indexed attestation and writes it to disk.
|
|
func (db *Store) SaveIndexedAttestation(idxAttestation *ethpb.IndexedAttestation) error {
|
|
key := encodeEpochSig(idxAttestation.Data.Target.Epoch, idxAttestation.Signature)
|
|
enc, err := proto.Marshal(idxAttestation)
|
|
if err != nil {
|
|
return errors.Wrap(err, "failed to marshal")
|
|
}
|
|
err = db.update(func(tx *bolt.Tx) error {
|
|
bucket := tx.Bucket(historicIndexedAttestationsBucket)
|
|
//if data is in db skip put and index functions
|
|
val := bucket.Get(key)
|
|
if val != nil {
|
|
return nil
|
|
}
|
|
createIndexedAttestationIndicesFromData(idxAttestation, tx)
|
|
if err := bucket.Put(key, enc); err != nil {
|
|
return errors.Wrap(err, "failed to include the indexed attestation in the historic indexed attestation bucket")
|
|
}
|
|
|
|
return err
|
|
})
|
|
|
|
// prune history to max size every PruneSlasherStoragePeriod epoch
|
|
if idxAttestation.Data.Source.Epoch%params.BeaconConfig().PruneSlasherStoragePeriod == 0 {
|
|
weakSubjectivityPeriod := params.BeaconConfig().WeakSubjectivityPeriod
|
|
err = db.PruneHistory(idxAttestation.Data.Source.Epoch, weakSubjectivityPeriod)
|
|
}
|
|
return err
|
|
}
|
|
|
|
func createIndexedAttestationIndicesFromData(idxAttestation *ethpb.IndexedAttestation, tx *bolt.Tx) error {
|
|
dataRoot, err := hashutil.HashProto(idxAttestation.Data)
|
|
|
|
if err != nil {
|
|
return errors.Wrap(err, "failed to hash indexed attestation data.")
|
|
}
|
|
protoIdxAtt := &slashpb.ValidatorIDToIdxAtt{
|
|
Signature: idxAttestation.Signature,
|
|
Indices: idxAttestation.AttestingIndices,
|
|
DataRoot: dataRoot[:],
|
|
}
|
|
key := bytesutil.Bytes8(idxAttestation.Data.Target.Epoch)
|
|
bucket := tx.Bucket(indexedAttestationsIndicesBucket)
|
|
enc := bucket.Get(key)
|
|
vIdxList, err := createValidatorIDsToIndexedAttestationList(enc)
|
|
if err != nil {
|
|
return errors.Wrap(err, "failed to decode value into ValidatorIDToIndexedAttestationList")
|
|
}
|
|
vIdxList.IndicesList = append(vIdxList.IndicesList, protoIdxAtt)
|
|
enc, err = proto.Marshal(vIdxList)
|
|
if err != nil {
|
|
return errors.Wrap(err, "failed to marshal")
|
|
}
|
|
if err := bucket.Put(key, enc); err != nil {
|
|
return errors.Wrap(err, "failed to include the indexed attestation in the historic indexed attestation bucket")
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// DeleteIndexedAttestation deletes a indexed attestation using the slot and its root as keys in their respective buckets.
|
|
func (db *Store) DeleteIndexedAttestation(idxAttestation *ethpb.IndexedAttestation) error {
|
|
key := encodeEpochSig(idxAttestation.Data.Target.Epoch, idxAttestation.Signature)
|
|
return db.update(func(tx *bolt.Tx) error {
|
|
bucket := tx.Bucket(historicIndexedAttestationsBucket)
|
|
enc := bucket.Get(key)
|
|
if enc == nil {
|
|
return nil
|
|
}
|
|
removeIndexedAttestationIndicesFromData(idxAttestation, tx)
|
|
if err := bucket.Delete(key); err != nil {
|
|
tx.Rollback()
|
|
return errors.Wrap(err, "failed to delete the indexed attestation from historic indexed attestation bucket")
|
|
}
|
|
return nil
|
|
})
|
|
}
|
|
|
|
func removeIndexedAttestationIndicesFromData(idxAttestation *ethpb.IndexedAttestation, tx *bolt.Tx) error {
|
|
dataRoot, err := hashutil.HashProto(idxAttestation.Data)
|
|
protoIdxAtt := &slashpb.ValidatorIDToIdxAtt{
|
|
Signature: idxAttestation.Signature,
|
|
Indices: idxAttestation.AttestingIndices,
|
|
DataRoot: dataRoot[:],
|
|
}
|
|
key := bytesutil.Bytes8(idxAttestation.Data.Target.Epoch)
|
|
bucket := tx.Bucket(indexedAttestationsIndicesBucket)
|
|
enc := bucket.Get(key)
|
|
vIdxList, err := createValidatorIDsToIndexedAttestationList(enc)
|
|
if err != nil {
|
|
return errors.Wrap(err, "failed to decode value into ValidatorIDToIndexedAttestationList")
|
|
}
|
|
for i, v := range vIdxList.IndicesList {
|
|
if reflect.DeepEqual(v, protoIdxAtt) {
|
|
copy(vIdxList.IndicesList[i:], vIdxList.IndicesList[i+1:])
|
|
vIdxList.IndicesList[len(vIdxList.IndicesList)-1] = nil // or the zero value of T
|
|
vIdxList.IndicesList = vIdxList.IndicesList[:len(vIdxList.IndicesList)-1]
|
|
break
|
|
}
|
|
}
|
|
enc, err = proto.Marshal(vIdxList)
|
|
if err != nil {
|
|
return errors.Wrap(err, "failed to marshal")
|
|
}
|
|
if err := bucket.Put(key, enc); err != nil {
|
|
return errors.Wrap(err, "failed to include the indexed attestation in the historic indexed attestation bucket")
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func (db *Store) pruneAttHistory(currentEpoch uint64, historySize uint64) error {
|
|
pruneTill := int64(currentEpoch) - int64(historySize)
|
|
if pruneTill <= 0 {
|
|
return nil
|
|
}
|
|
return db.update(func(tx *bolt.Tx) error {
|
|
bucket := tx.Bucket(historicIndexedAttestationsBucket)
|
|
c := tx.Bucket(historicIndexedAttestationsBucket).Cursor()
|
|
max := bytesutil.Bytes8(uint64(pruneTill))
|
|
for k, _ := c.First(); k != nil && bytes.Compare(k[:8], max) <= 0; k, _ = c.Next() {
|
|
if err := bucket.Delete(k); err != nil {
|
|
return errors.Wrap(err, "failed to delete the indexed attestation from historic indexed attestation bucket")
|
|
}
|
|
}
|
|
idxBucket := tx.Bucket(indexedAttestationsIndicesBucket)
|
|
c = tx.Bucket(indexedAttestationsIndicesBucket).Cursor()
|
|
for k, _ := c.First(); k != nil && bytes.Compare(k[:8], max) <= 0; k, _ = c.Next() {
|
|
if err := idxBucket.Delete(k); err != nil {
|
|
return errors.Wrap(err, "failed to delete the indexed attestation from indexed attestation indexes bucket")
|
|
}
|
|
}
|
|
return nil
|
|
})
|
|
}
|