package parlia import ( "bytes" "context" "encoding/binary" "encoding/hex" "errors" "fmt" "io" "math" "math/big" "math/rand" "sort" "strings" "sync" "time" "github.com/bits-and-blooms/bitset" "github.com/ethereum/go-ethereum/common/lru" "github.com/ethereum/go-ethereum/common/pepper8" "github.com/ethereum/go-ethereum/common/pipe8" "github.com/holiman/uint256" "github.com/prysmaticlabs/prysm/v5/crypto/bls" "golang.org/x/crypto/sha3" "github.com/ethereum/go-ethereum/accounts" "github.com/ethereum/go-ethereum/accounts/abi" "github.com/ethereum/go-ethereum/common" "github.com/ethereum/go-ethereum/common/gopool" "github.com/ethereum/go-ethereum/common/hexutil" cmath "github.com/ethereum/go-ethereum/common/math" "github.com/ethereum/go-ethereum/common/systemcontract" "github.com/ethereum/go-ethereum/consensus" "github.com/ethereum/go-ethereum/consensus/misc/eip1559" "github.com/ethereum/go-ethereum/consensus/misc/eip4844" "github.com/ethereum/go-ethereum/core" "github.com/ethereum/go-ethereum/core/forkid" "github.com/ethereum/go-ethereum/core/state" "github.com/ethereum/go-ethereum/core/systemcontracts" "github.com/ethereum/go-ethereum/core/tracing" "github.com/ethereum/go-ethereum/core/types" "github.com/ethereum/go-ethereum/core/vm" "github.com/ethereum/go-ethereum/crypto" "github.com/ethereum/go-ethereum/ethdb" "github.com/ethereum/go-ethereum/internal/ethapi" "github.com/ethereum/go-ethereum/log" "github.com/ethereum/go-ethereum/metrics" "github.com/ethereum/go-ethereum/params" "github.com/ethereum/go-ethereum/rlp" "github.com/ethereum/go-ethereum/rpc" "github.com/ethereum/go-ethereum/trie" ) const ( inMemorySnapshots = 1280 // Number of recent snapshots to keep in memory; a buffer exceeding the EpochLength inMemorySignatures = 4096 // Number of recent block signatures to keep in memory inMemoryFrequencies = 128 // Number of post-Snake8Fix frequency-data entries to memoize by parent hash inMemoryHeaders = 86400 // Number of recent headers to keep in memory for double sign detection, checkpointInterval = 1024 // Number of blocks after which to save the snapshot to the database lorentzEpochLength uint64 = 500 // Epoch length starting from the Lorentz hard fork maxwellEpochLength uint64 = 1000 // Epoch length starting from the Maxwell hard fork defaultBlockInterval uint64 = 3000 // Default block interval in milliseconds lorentzBlockInterval uint64 = 1500 // Block interval starting from the Lorentz hard fork maxwellBlockInterval uint64 = 750 // Block interval starting from the Maxwell hard fork fermiBlockInterval uint64 = 450 // Block interval starting from the Fermi hard fork defaultTurnLength uint8 = 1 // Default consecutive number of blocks a validator receives priority for block production snake8TurnLength uint8 = 50 // Chiliz Snake8 forces this turn length whenever the fork is active extraVanity = 32 // Fixed number of extra-data prefix bytes reserved for signer vanity extraSeal = 65 // Fixed number of extra-data suffix bytes reserved for signer seal nextForkHashSize = 4 // Fixed number of extra-data suffix bytes reserved for nextForkHash. turnLengthSize = 1 // Fixed number of extra-data suffix bytes reserved for turnLength validatorBytesLengthBeforeLuban = common.AddressLength validatorBytesLength = common.AddressLength + types.BLSPublicKeyLength validatorNumberSize = 1 // Fixed number of extra prefix bytes reserved for validator number after Luban wiggleTime uint64 = 1000 // milliseconds, Random delay (per signer) to allow concurrent signers defaultInitialBackOffTime uint64 = 1000 // milliseconds, Default backoff time for the second validator permitted to produce blocks lorentzInitialBackOffTime uint64 = 2000 // milliseconds, Backoff time for the second validator permitted to produce blocks from the Lorentz hard fork systemRewardPercent = 5 // it means 1/3 percentage of gas fee incoming will be distributed to system collectAdditionalVotesRewardRatio = 100 // ratio of additional reward for collecting more votes than needed, the denominator is 100 gasLimitBoundDivisorBeforeLorentz uint64 = 256 // The bound divisor of the gas limit, used in update calculations before lorentz hard fork. // `finalityRewardInterval` should be smaller than `inMemorySnapshots`, otherwise, it will result in excessive computation. finalityRewardInterval = 200 kAncestorGenerationDepth = 3 ) var ( defaultEpochLength = uint64(200) // Default number of blocks of checkpoint to update validatorSet from contract diffInTurn = big.NewInt(2) // Block difficulty for in-turn signatures diffNoTurn = big.NewInt(1) // Block difficulty for out-of-turn signatures validatorFrequencyDataPrefix = []byte("VFQ") // Prefix for validator frequency data in the header.Extra field // 100 native token maxSystemBalance = new(uint256.Int).Mul(uint256.NewInt(100), uint256.NewInt(params.Ether)) verifyVoteAttestationErrorCounter = metrics.NewRegisteredCounter("parlia/verifyVoteAttestation/error", nil) updateAttestationErrorCounter = metrics.NewRegisteredCounter("parlia/updateAttestation/error", nil) validVotesfromSelfCounter = metrics.NewRegisteredCounter("parlia/VerifyVote/self", nil) doubleSignCounter = metrics.NewRegisteredCounter("parlia/doublesign", nil) snake8FreqVerifySkippedCounter = metrics.NewRegisteredCounter("parlia/snake8fix/verifySkipped", nil) intentionalDelayMiningCounter = metrics.NewRegisteredCounter("parlia/intentionalDelayMining", nil) attestationVoteCountGauge = metrics.NewRegisteredGauge("parlia/attestation/voteCount", nil) ) // Various error messages to mark blocks invalid. These should be private to // prevent engine specific errors from being referenced in the remainder of the // codebase, inherently breaking if the engine is swapped out. Please put common // error types into the consensus package. var ( // errUnknownBlock is returned when the list of validators is requested for a block // that is not part of the local blockchain. errUnknownBlock = errors.New("unknown block") // errMissingVanity is returned if a block's extra-data section is shorter than // 32 bytes, which is required to store the signer vanity. errMissingVanity = errors.New("extra-data 32 byte vanity prefix missing") // errMissingSignature is returned if a block's extra-data section doesn't seem // to contain a 65 byte secp256k1 signature. errMissingSignature = errors.New("extra-data 65 byte signature suffix missing") // errExtraValidators is returned if non-sprint-end block contain validator data in // their extra-data fields. errExtraValidators = errors.New("non-sprint-end block contains extra validator list") // errInvalidSpanValidators is returned if a block contains an // invalid list of validators (i.e. non divisible by 20 bytes). errInvalidSpanValidators = errors.New("invalid validator list on sprint end block") // errInvalidTurnLength is returned if a block contains an // invalid length of turn (i.e. no data left after parsing validators). errInvalidTurnLength = errors.New("invalid turnLength") // errInvalidSnake8Extra is returned if a Snake8 block's extra-data does not // contain a well-formed validator-frequency block (missing/short/misplaced // prefix or embedded parent timestamp). errInvalidSnake8Extra = errors.New("invalid snake8 frequency extra data") // errMismatchedSnake8ParentTime is returned if the parent timestamp embedded // in a Snake8 block's extra-data does not match the real parent.Time. This // guards against forged activation timestamps causing a consensus split. errMismatchedSnake8ParentTime = errors.New("snake8 embedded parent timestamp mismatch") // errMismatchedSnake8FrequencyData is returned (post-Snake8Fix, COR-173) if the // validator-frequency bytes embedded in a header do not match the recomputation // from the parent block's state. Pre-Snake8Fix the embedded bytes are // producer-controlled and unverified (historical behavior). errMismatchedSnake8FrequencyData = errors.New("snake8 embedded frequency data mismatch") // errSnake8StakeLookup marks a NODE-LOCAL failure to read validator stakes for // the Snake8 frequency data (RPC timeout under load, gas cap, missing historical // state during tracing replay, engine built without an RPC backend). It says // nothing about the block's validity: sealing treats it as fatal, verification // fails open on it (verifySnake8FrequencyData). errSnake8StakeLookup = errors.New("snake8fix stake lookup failed") // errInvalidMixDigest is returned if a block's mix digest is non-zero. errInvalidMixDigest = errors.New("non-zero mix digest") // errInvalidUncleHash is returned if a block contains an non-empty uncle list. errInvalidUncleHash = errors.New("non empty uncle hash") // errMismatchingEpochValidators is returned if a sprint block contains a // list of validators different than the one the local node calculated. errMismatchingEpochValidators = errors.New("mismatching validator list on epoch block") // errMismatchingEpochTurnLength is returned if a sprint block contains a // turn length different than the one the local node calculated. errMismatchingEpochTurnLength = errors.New("mismatching turn length on epoch block") // errInvalidDifficulty is returned if the difficulty of a block is missing. errInvalidDifficulty = errors.New("invalid difficulty") // errWrongDifficulty is returned if the difficulty of a block doesn't match the // turn of the signer. errWrongDifficulty = errors.New("wrong difficulty") // errOutOfRangeChain is returned if an authorization list is attempted to // be modified via out-of-range or non-contiguous headers. errOutOfRangeChain = errors.New("out of range or non-contiguous chain") // errBlockHashInconsistent is returned if an authorization list is attempted to // insert an inconsistent block. errBlockHashInconsistent = errors.New("the block hash is inconsistent") // errUnauthorizedValidator is returned if a header is signed by a non-authorized entity. errUnauthorizedValidator = func(val string) error { return errors.New("unauthorized validator: " + val) } // errCoinBaseMisMatch is returned if a header's coinbase do not match with signature errCoinBaseMisMatch = errors.New("coinbase do not match with signature") // errRecentlySigned is returned if a header is signed by an authorized entity // that already signed a header recently, thus is temporarily not allowed to. errRecentlySigned = errors.New("recently signed") ) // SignerFn is a signer callback function to request a header to be signed by a // backing account. type SignerFn func(accounts.Account, string, []byte) ([]byte, error) type SignerTxFn func(accounts.Account, *types.Transaction, *big.Int) (*types.Transaction, error) func isToSystemContract(to common.Address) bool { return systemcontract.IsSystemContract(to) || to == pepper8.Pepper8RecipientAddress || to == pipe8.Pipe8RecipientAddress } // ecrecover extracts the Ethereum account address from a signed header. func ecrecover(header *types.Header, sigCache *lru.Cache[common.Hash, common.Address], chainId *big.Int) (common.Address, error) { // If the signature's already cached, return that hash := header.Hash() if address, known := sigCache.Get(hash); known { return address, nil } // Retrieve the signature from the header extra-data if len(header.Extra) < extraSeal { return common.Address{}, errMissingSignature } signature := header.Extra[len(header.Extra)-extraSeal:] // Recover the public key and the Ethereum address pubkey, err := crypto.Ecrecover(types.SealHash(header, chainId).Bytes(), signature) if err != nil { return common.Address{}, err } var signer common.Address copy(signer[:], crypto.Keccak256(pubkey[1:])[12:]) sigCache.Add(hash, signer) return signer, nil } // ParliaRLP returns the rlp bytes which needs to be signed for the parlia // sealing. The RLP to sign consists of the entire header apart from the 65 byte signature // contained at the end of the extra data. // // Note, the method requires the extra data to be at least 65 bytes, otherwise it // panics. This is done to avoid accidentally using both forms (signature present // or not), which could be abused to produce different hashes for the same header. func ParliaRLP(header *types.Header, chainId *big.Int) []byte { b := new(bytes.Buffer) types.EncodeSigHeader(b, header, chainId) return b.Bytes() } // Parlia is the consensus engine of BSC type Parlia struct { chainConfig *params.ChainConfig // Chain config config *params.ParliaConfig // Consensus engine configuration parameters for parlia consensus genesisHash common.Hash db ethdb.Database // Database to store and retrieve snapshot checkpoints recentSnaps *lru.Cache[common.Hash, *Snapshot] // Snapshots for recent block to speed up signatures *lru.Cache[common.Hash, common.Address] // Signatures of recent blocks to speed up mining recentHeaders *lru.Cache[string, common.Hash] // Recent headers to check for double signing: key includes block number and miner. value is the block header // If same key's value already exists for different block header roots then double sign is detected signer types.Signer val common.Address // Ethereum address of the signing key signFn SignerFn // Signer function to authorize hashes with signTxFn SignerTxFn lock sync.RWMutex // Protects the signer fields ethAPI *ethapi.BlockChainAPI VotePool consensus.VotePool validatorSetABIBeforeLuban abi.ABI validatorSetABI abi.ABI slashABI abi.ABI tokenomicsABI abi.ABI stakeHubABI abi.ABI // stakeReader fetches a validator's total delegated stake for the Snake8 // frequency data. Defaults to getValidatorTotalDelegated (refreshFrequencyRLP // also falls back to it when nil, so struct-literal engines don't trap); tests // substitute a deterministic source since engines built without an RPC backend // cannot perform contract calls. stakeReader func(validator common.Address, blockNumber uint64, state *rpc.BlockNumberOrHash) (*big.Int, error) // frequencyCache memoizes post-Snake8Fix frequency bytes by parent hash — see // refreshFrequencyRLP for why this is sound and why pre-fork bytes never enter it. frequencyCache *lru.Cache[common.Hash, []byte] } // New creates a Parlia consensus engine. func New( chainConfig *params.ChainConfig, db ethdb.Database, ethAPI *ethapi.BlockChainAPI, genesisHash common.Hash, ) *Parlia { // get parlia config parliaConfig := chainConfig.Parlia log.Info("Parlia", "chainConfig", chainConfig) // Override default epoch length with value from genesis config defaultEpochLength = parliaConfig.Epoch vABIBeforeLuban, err := abi.JSON(strings.NewReader(validatorSetABIBeforeLuban)) if err != nil { panic(err) } vABI, err := abi.JSON(strings.NewReader(validatorSetABI)) if err != nil { panic(err) } // signableSystemTxSelectors exists only for BEP-675 bid blocks, and covers BSC // validator-set methods that Chiliz's BAS validator set does not implement. // Assert the selectors only on chains that schedule the forks which pack them. if chainConfig.PlatoBlock != nil && chainConfig.FeynmanTime != nil { for methodName, selector := range signableSystemTxSelectors { method, ok := vABI.Methods[methodName] if !ok || !bytes.Equal(method.ID, selector[:]) { panic(fmt.Sprintf("invalid validator set ABI selector for %s", methodName)) } } } sABI, err := abi.JSON(strings.NewReader(slashABI)) if err != nil { panic(err) } tABI, err := abi.JSON(strings.NewReader(tokenomicsABI)) if err != nil { panic(err) } stABI, err := abi.JSON(strings.NewReader(stakeABI)) if err != nil { panic(err) } c := &Parlia{ chainConfig: chainConfig, config: parliaConfig, genesisHash: genesisHash, db: db, ethAPI: ethAPI, recentSnaps: lru.NewCache[common.Hash, *Snapshot](inMemorySnapshots), recentHeaders: lru.NewCache[string, common.Hash](inMemoryHeaders), frequencyCache: lru.NewCache[common.Hash, []byte](inMemoryFrequencies), signatures: lru.NewCache[common.Hash, common.Address](inMemorySignatures), validatorSetABIBeforeLuban: vABIBeforeLuban, validatorSetABI: vABI, slashABI: sABI, tokenomicsABI: tABI, stakeHubABI: stABI, signer: types.LatestSigner(chainConfig), } c.stakeReader = c.getValidatorTotalDelegated return c } func (p *Parlia) IsSystemTransaction(tx *types.Transaction, header *types.Header) (bool, error) { if tx.To() == nil || !isToSystemContract(*tx.To()) { return false, nil } if tx.EffectiveGasPriceForBSC().Sign() != 0 { return false, nil } sender, err := types.Sender(p.signer, tx) if err != nil { return false, errors.New("UnAuthorized transaction") } return sender == header.Coinbase, nil } // tokenomicsDepositSelector is the 4-byte ABI selector of // Tokenomics.deposit(address,uint256,uint256) (keccak256 prefix 0x0efe6a8b). var tokenomicsDepositSelector = [4]byte{0x0e, 0xfe, 0x6a, 0x8b} // IsTokenomicsDeposit returns true if to address is the tokenomics contract and tx data // starts with Tokenomics.deposit() method signature. // // to must be non-nil: every caller reaches this only for a transaction that // IsSystemTransaction / IsSystemContract already accepted, both of which reject // a nil destination. The guard enforces that rather than trusting it: CLAUDE.md // section 5 anticipates new replay and simulation paths calling in here, and a // caller that forgets it would panic the node on a contract-creation tx instead // of getting false. func (p *Parlia) IsTokenomicsDeposit(to *common.Address, data []byte) bool { return to != nil && *to == systemcontract.TokenomicsContractAddress && len(data) >= len(tokenomicsDepositSelector) && bytes.Equal(data[:len(tokenomicsDepositSelector)], tokenomicsDepositSelector[:]) } // IsPepper8Deposit returns true if to address is the pepper8 recipient and from // address is coinbase. Nil-total, for the reason on IsTokenomicsDeposit. func (p *Parlia) IsPepper8Deposit(from *common.Address, to *common.Address, coinbase *common.Address) bool { if from == nil || to == nil || coinbase == nil { return false } isDestinationPepper8Recipient := bytes.Equal(to.Bytes(), pepper8.Pepper8RecipientAddress.Bytes()) isFromCoinbase := bytes.Equal(from.Bytes(), coinbase.Bytes()) return isDestinationPepper8Recipient && isFromCoinbase } // IsPipe8Deposit returns true if to address is the pipe8 recipient and from // address is coinbase. Nil-total, for the reason on IsTokenomicsDeposit. func (p *Parlia) IsPipe8Deposit(from *common.Address, to *common.Address, coinbase common.Address) bool { if from == nil || to == nil { return false } isDestinationPipe8Recipient := bytes.Equal(to.Bytes(), pipe8.Pipe8RecipientAddress.Bytes()) isFromCoinbase := bytes.Equal(from.Bytes(), coinbase.Bytes()) return isDestinationPipe8Recipient && isFromCoinbase } func (p *Parlia) IsSystemContract(to *common.Address) bool { if to == nil { return false } return isToSystemContract(*to) } // Author implements consensus.Engine, returning the SystemAddress func (p *Parlia) Author(header *types.Header) (common.Address, error) { return header.Coinbase, nil } // ConsensusAddress returns the consensus address of the validator func (p *Parlia) ConsensusAddress() common.Address { return p.val } // VerifyHeader checks whether a header conforms to the consensus rules. func (p *Parlia) VerifyHeader(chain consensus.ChainHeaderReader, header *types.Header) error { return p.verifyHeader(chain, header, nil) } // VerifyHeaders is similar to VerifyHeader, but verifies a batch of headers. The // method returns a quit channel to abort the operations and a results channel to // retrieve the async verifications (the order is that of the input slice). func (p *Parlia) VerifyHeaders(chain consensus.ChainHeaderReader, headers []*types.Header) (chan<- struct{}, <-chan error) { abort := make(chan struct{}) results := make(chan error, len(headers)) gopool.Submit(func() { for i, header := range headers { err := p.verifyHeader(chain, header, headers[:i]) select { case <-abort: return case results <- err: } } }) return abort, results } // getValidatorBytesFromHeader returns the validators bytes extracted from the header's extra field if exists. // The validators bytes would be contained only in the epoch block's header, and its each validator bytes length is fixed. // On luban fork, we introduce vote attestation into the header's extra field, so extra format is different from before. // Before luban fork: |---Extra Vanity---|---Validators Bytes (or Empty)---|---Extra Seal---| // After luban fork: |---Extra Vanity---|---Validators Number and Validators Bytes (or Empty)---|---Vote Attestation (or Empty)---|---Extra Seal---| // After bohr fork: |---Extra Vanity---|---Validators Number and Validators Bytes (or Empty)---|---Turn Length (or Empty)---|---Vote Attestation (or Empty)---|---Extra Seal---| // After snake8 fork: |---Extra Vanity---|---Validators Bytes (or Empty) ---|---Turn Length (or Empty)---/---Vote Attestation (or Empty)---/---Frequency Data Prefix---|---Parent Timestamp---|---Frequency data---|---Extra Seal---| func getValidatorBytesFromHeader(header *types.Header, chainConfig *params.ChainConfig, epochLength uint64) []byte { if len(header.Extra) <= extraVanity+extraSeal { return nil } if !chainConfig.IsLuban(header.Number) { start := extraVanity end := len(header.Extra) - extraSeal body := header.Extra[start:end] // Legacy quirks, preserved byte-for-byte on BOTH epoch and non-epoch // headers. The pre-COR-213 parser rejected a header outright when // // a) the body began with "VFQ", or // b) the first "VFQ" occurrence in header.Extra began inside the vanity, // including one straddling the vanity end into the first validator // address — the scan started at offset 0, so such an occurrence set // end < start and fell into the `end <= start` return. // // Both are consensus rules on the live networks, not artefacts of the // scan: relaxing either accepts an epoch header that every 2.10.6 node // rejects, which splits the chain across a rolling upgrade instead of // halting it uniformly. They are therefore checked ahead of the // epoch/non-epoch split rather than inside the non-epoch branch. // // The residual cost is narrow and deliberate: an epoch header whose // lowest-sorted validator address itself begins with 0x564651 stays // unparseable. Legacy rejected it too, so this is not a regression — it is // the one corner of COR-213 that cannot be fixed without a hardfork. // // The legacy loop ran i over [0, end-3], i.e. it never read prefix bytes // from the seal; the straddle window is therefore bounded by end. legacyScanEnd := start + len(validatorFrequencyDataPrefix) - 1 if legacyScanEnd > end { legacyScanEnd = end } if bytes.HasPrefix(body, validatorFrequencyDataPrefix) || bytes.Contains(header.Extra[:legacyScanEnd], validatorFrequencyDataPrefix) { return nil } if header.Number.Uint64()%epochLength != 0 { // Non-epoch headers carry no validator list: the body is either the // frequency block (caught by the HasPrefix guard above) or something // verifyHeader rejects as errExtraValidators. return body } // Epoch header: |validators (n*20 bytes)|"VFQ"|parent ts|frequency RLP| // once Snake8 is active, |validators| before. The frequency block is // located structurally (validator-aligned prefix + well-formed tail), never // by a byte-wise scan for "VFQ": an address containing those three bytes // used to truncate the list (COR-213). switch off, found, hostile := preLubanFrequencyBlockOffset(body); { case hostile: return nil case found && off == 0: // The frequency block starts immediately: no validator entries. return nil case found: end = start + off } if (end-start)%validatorBytesLengthBeforeLuban != 0 { return nil } return header.Extra[start:end] } if header.Number.Uint64()%epochLength != 0 { return nil } num := int(header.Extra[extraVanity]) start := extraVanity + validatorNumberSize end := start + num*validatorBytesLength extraMinLen := end + extraSeal if chainConfig.IsBohr(header.Number, header.Time) { extraMinLen += turnLengthSize } if num == 0 || len(header.Extra) < extraMinLen { return nil } return header.Extra[start:end] } // frequencyBlockEntry is the wire shape of one Snake8 frequency-data entry: // what calcFrequencyRLP encodes and selectValidatorFromFrequencyRLP decodes. type frequencyBlockEntry struct { Address common.Address Frequency *big.Int } // maxFrequencyBlockCandidates bounds the structural decodes attempted while // locating the frequency block in one header, so a hostile header cannot make // the parser do quadratic work. An honest epoch header has exactly one // candidate; one candidate per validator entry is already absurd. const maxFrequencyBlockCandidates = 256 // decodesAsFrequencyList reports whether tail is exactly one NON-EMPTY RLP list // of (address, frequency) entries, the encoding calcFrequencyRLP produces. // Trailing bytes, a different structure, or an empty input all fail. // // The empty list (0xc0) is rejected on purpose: calcFrequencyRLP returns an // error rather than an empty list when there are no eligible candidates, so the // producer emits either a list with at least one entry or nothing at all — an // empty list is not a shape any honest header carries. Accepting it would make // a second decodable candidate cost an attacker a single byte: with "VFQ" at // bytes 11..13 of their own (freely ground) address, a frequency whose encoding // ends in 0xc0 and a list length ≡ 1 (mod 20), the last twelve bytes of the // honest frequency list form a 20-aligned candidate with tail 0xc0. Two // decodable candidates are reported hostile, getValidatorBytesFromHeader // returns nil, and every node rejects the epoch block — the COR-213 halt the // structural search exists to prevent. Requiring an entry raises the bar to 25 // attacker-controlled bytes at the very end of the body, which the producer's // encoding never places there. See TestPreLubanFrequencyEmptyListDecoy. func decodesAsFrequencyList(tail []byte) bool { if len(tail) == 0 { return false } var entries []frequencyBlockEntry return rlp.DecodeBytes(tail, &entries) == nil && len(entries) > 0 } // preLubanFrequencyBlockOffset locates the Snake8 frequency block inside the // body (vanity and seal stripped) of a pre-Luban epoch header, where it follows // the fixed-width 20-byte validator entries: // // |validator 1|...|validator n|"VFQ"|LE parent timestamp|frequency RLP| // // The prefix alone is not a delimiter: a validator address may contain, or even // start with, the bytes 0x564651, and the previous byte-wise scan for their // first occurrence truncated the validator list there (COR-213). The block is // therefore located structurally. A candidate is an offset that is a multiple // of the validator width (zero included: a first address may itself start with // the prefix), starts with "VFQ" and has room for the 8-byte timestamp. Its tail // is everything up to the seal, and refreshFrequencyRLP only ever embeds two // tail shapes: the RLP list of entries, or nothing at all when calcFrequencyRLP // failed. So: // // - exactly one candidate whose tail decodesAsFrequencyList: that is the // block. The bytes after a validator address are the remaining raw // 20-byte entries, which never form such a list, so an address is told // apart from the block in every realistic layout; // - more than one such candidate: unreachable for an honest producer and // not constructible by an attacker who controls only their own address. // Reported as hostile and rejected; // - none, but a candidate whose tail is empty (necessarily the last offset): // the degraded block with no frequency data. It only counts when no // decodable candidate exists, because an attacker can cheaply forge one: // their address appears inside the honest frequency list, and placing // "VFQ" at the right position in it puts the prefix exactly eleven bytes // before the seal, on an aligned offset. FuzzHeaderExtraRoundTrip found // that shape (see its testdata) against an earlier rule that let it win; // - none at all: a pre-Snake8 header, the whole body is validator entries. // // Only a bounded number of prefixed offsets is examined, so a hostile header // cannot make the parser do quadratic work; exceeding the bound is hostile. func preLubanFrequencyBlockOffset(body []byte) (offset int, found bool, hostile bool) { blockHeader := len(validatorFrequencyDataPrefix) + 8 // "VFQ" + parent timestamp attempts, decodable, emptyAtEnd := 0, 0, -1 for off := 0; off+blockHeader <= len(body); off += validatorBytesLengthBeforeLuban { if !bytes.HasPrefix(body[off:], validatorFrequencyDataPrefix) { continue } tail := body[off+blockHeader:] if len(tail) == 0 { emptyAtEnd = off continue } attempts++ if attempts > maxFrequencyBlockCandidates { return 0, false, true } if decodesAsFrequencyList(tail) { offset = off decodable++ if decodable > 1 { // The verdict can no longer change; stop decoding. Identical // result to falling through to the switch below, but it caps // the RLP work a hostile header can provoke at two decodes. return 0, false, true } } } switch { case decodable == 1: return offset, true, false case decodable > 1: return 0, false, true case emptyAtEnd >= 0: return emptyAtEnd, true, false } return 0, false, false } // getVoteAttestationFromHeader returns the vote attestation extracted from the header's extra field if exists. func getVoteAttestationFromHeader(header *types.Header, chainConfig *params.ChainConfig, epochLength uint64) (*types.VoteAttestation, error) { if len(header.Extra) <= extraVanity+extraSeal { return nil, nil } if !chainConfig.IsLuban(header.Number) { return nil, nil } var attestationBytes []byte if header.Number.Uint64()%epochLength != 0 { attestationBytes = header.Extra[extraVanity : len(header.Extra)-extraSeal] } else { num := int(header.Extra[extraVanity]) start := extraVanity + validatorNumberSize + num*validatorBytesLength if chainConfig.IsBohr(header.Number, header.Time) { start += turnLengthSize } end := len(header.Extra) - extraSeal if end <= start { return nil, nil } attestationBytes = header.Extra[start:end] } var attestation types.VoteAttestation if err := rlp.Decode(bytes.NewReader(attestationBytes), &attestation); err != nil { return nil, fmt.Errorf("block %d has vote attestation info, decode err: %s", header.Number.Uint64(), err) } return &attestation, nil } // getParent returns the parent of a given block. func (p *Parlia) getParent(chain consensus.ChainHeaderReader, header *types.Header, parents []*types.Header) (*types.Header, error) { var parent *types.Header number := header.Number.Uint64() if len(parents) > 0 { parent = parents[len(parents)-1] } else { parent = chain.GetHeader(header.ParentHash, number-1) } if parent == nil || parent.Number.Uint64() != number-1 || parent.Hash() != header.ParentHash { return nil, consensus.ErrUnknownAncestor } return parent, nil } // isSnake8Enabled returns true if the parent block's timestamp >= snake8Time. // // When the parent header is available locally (the common case) the decision is // made from the trusted parent.Time. During batch sync the parent may not yet be // in the local store, in which case we provisionally fall back to the parent // timestamp embedded in header.Extra. That embedded value is producer-controlled // and MUST NOT be trusted on its own; verifySnake8Extra cross-validates it // against the real parent.Time during full verification, so a forged value can // never make a block accepted on the wrong side of the fork boundary. func (p *Parlia) isSnake8Enabled(chain consensus.ChainHeaderReader, header *types.Header) bool { parent := chain.GetHeaderByHash(header.ParentHash) if parent != nil { return p.chainConfig.IsSnake8(parent.Time) } // Parent not in cache: provisionally read the embedded parent timestamp using // the strict, bounds-checked extractor. Any malformed layout yields false // rather than an out-of-bounds read on attacker-controlled bytes. ts, ok := extractSnake8ParentTimestamp(header, p.chainConfig) if !ok { log.Trace("failed to extract parent timestamp from extra", "number", header.Number.Uint64()) return false } return p.chainConfig.IsSnake8(ts) } // verifySnake8Extra authoritatively validates the Snake8 fork data embedded in a // header's extra-data against the real (now-known) parent header. Fork // activation is decided solely from the trusted parent.Time, never from // producer-controlled bytes: // // - When Snake8 is active for this block, the header must carry a well-formed // validator-frequency block whose embedded parent timestamp equals // parent.Time. A forged or absent timestamp is rejected here, so a malicious // producer cannot push syncing nodes onto the wrong side of the fork. // - When Snake8 is not yet active, the header must not carry a frequency block. func (p *Parlia) verifySnake8Extra(header, parent *types.Header) error { ts, ok := extractSnake8ParentTimestamp(header, p.chainConfig) if !p.chainConfig.IsSnake8(parent.Time) { // Pre-fork block: it must not carry a Snake8 frequency block. if ok { return errInvalidSnake8Extra } return nil } // Snake8 is active: a well-formed frequency block is mandatory and its // embedded parent timestamp must match the real parent. if !ok { return errInvalidSnake8Extra } if ts != parent.Time { return errMismatchedSnake8ParentTime } return nil } // verifySnake8FrequencyData authoritatively validates the validator-frequency // bytes embedded in a Snake8 header against their recomputation from the parent // block's state (COR-173, Snake8Fix fork). Pre-Snake8Fix the bytes are accepted // as-is: historical blocks were produced with unpinned latest-state stake reads // (and soft degradation to zero stakes on lookup failure), so a strict // recomputation could reject valid history — the check is therefore gated on the // parent's timestamp, mirroring Snake8's own activation rule. // // The check is state-dependent, so it runs from Finalize (like verifyValidators), // not from VerifyHeader: a node executing block N necessarily has state N-1. // refreshFrequencyRLP is deterministic post-fork (stakes pinned to the parent // block, candidates from the parent snapshot's Recents), so every honest node // recomputes exactly the bytes an honest producer embedded. // // Error classification: only a byte MISMATCH condemns the block. A node-local // inability to read stakes (errSnake8StakeLookup: RPC timeout under load, gas cap, // no historical state during tracing replay, nil-ethAPI engines like `geth import`) // fails OPEN — logged loudly and counted — because an error out of Finalize is // treated as a bad block by insertChain: failing closed would let a local condition // (or RPC-load DoS) make a node reject the honest head and stall. Fail-open only // relaxes the check on nodes that could not have evaluated it; sealing keeps the // hard failure, so an honest producer never emits unverifiable blocks. func (p *Parlia) verifySnake8FrequencyData(chain consensus.ChainHeaderReader, header, parent *types.Header) error { if !p.chainConfig.IsSnake8Fix(parent.Time) { return nil } embedded, err := parseValidatorFrequencies(header, p.chainConfig) if err != nil { return err } number := header.Number.Uint64() // isSnake8Fork=true so snapshot() returns a caller-owned copy (COR-174) that // refreshFrequencyRLP may write to; extraHeader=nil because the embedded // bytes under judgment must not pre-stamp the snapshot we recompute into. snap, err := p.snapshot(chain, number-1, header.ParentHash, nil, true, nil) if err != nil { return err } if err := p.refreshFrequencyRLP(snap, number, parent); err != nil { if errors.Is(err, errSnake8StakeLookup) { snake8FreqVerifySkippedCounter.Inc(1) log.Error("Skipping Snake8 frequency verification: cannot read stakes locally", "number", number, "hash", header.Hash(), "err", err) return nil } return err } if !bytes.Equal(embedded, snap.FrequencyRLP) { log.Warn("Snake8 embedded frequency data does not match state recomputation", "number", number, "hash", header.Hash(), "coinbase", header.Coinbase, "embedded", hex.EncodeToString(embedded), "recomputed", hex.EncodeToString(snap.FrequencyRLP)) return errMismatchedSnake8FrequencyData } return nil } // trimParents safely removes last element if exists. func trimParents(parents []*types.Header) []*types.Header { if len(parents) > 1 { return parents[:len(parents)-1] } return nil } // verifyVoteAttestation checks whether the vote attestation in the header is valid. func (p *Parlia) verifyVoteAttestation(chain consensus.ChainHeaderReader, header *types.Header, parents []*types.Header) error { // === Step 1: Extract attestation === epochLength, err := p.epochLength(chain, header, parents) if err != nil { return err } attestation, err := getVoteAttestationFromHeader(header, chain.Config(), epochLength) if err != nil { return err } if attestation == nil { return nil } if attestation.Data == nil { return errors.New("invalid attestation, vote data is nil") } if len(attestation.Extra) > types.MaxAttestationExtraLength { return fmt.Errorf("invalid attestation, too large extra length: %d", len(attestation.Extra)) } if attestation.Data.SourceNumber >= attestation.Data.TargetNumber { return errors.New("invalid attestation, SourceNumber not lower than TargetNumber") } // === Step 2: Verify source block === parent, err := p.getParent(chain, header, parents) if err != nil { return err } // The source block should be the highest justified block. sourceNumber := attestation.Data.SourceNumber sourceHash := attestation.Data.SourceHash headers := []*types.Header{parent} if len(parents) > 0 { headers = parents } justifiedBlockNumber, justifiedBlockHash, err := p.GetJustifiedNumberAndHash(chain, headers) if err != nil { return errors.New("unexpected error when getting the highest justified number and hash") } if sourceNumber != justifiedBlockNumber || sourceHash != justifiedBlockHash { return fmt.Errorf("invalid attestation, source mismatch, expected block: %d, hash: %s; real block: %d, hash: %s", justifiedBlockNumber, justifiedBlockHash, sourceNumber, sourceHash) } // === Step 3: Verify target block === targetNumber := attestation.Data.TargetNumber targetHash := attestation.Data.TargetHash match := false ancestor := parent ancestorParents := trimParents(parents) for range p.GetAncestorGenerationDepth(header) { if targetNumber == ancestor.Number.Uint64() && targetHash == ancestor.Hash() { match = true break } ancestor, err = p.getParent(chain, ancestor, ancestorParents) if err != nil { return err } ancestorParents = trimParents(ancestorParents) } if !match { return fmt.Errorf("invalid attestation, target mismatch, real block: %d, hash: %s", targetNumber, targetHash) } // === Step 4: Check quorum === // The snapshot should be the targetNumber-1 block's snapshot. snap, err := p.snapshot(chain, ancestor.Number.Uint64()-1, ancestor.ParentHash, ancestorParents, p.isSnake8Enabled(chain, ancestor), ancestor) if err != nil { return err } // Filter out valid validator from attestation. validators := snap.validators() validatorsBitSet := bitset.From([]uint64{uint64(attestation.VoteAddressSet)}) if validatorsBitSet.Count() > uint(len(validators)) { return errors.New("invalid attestation, vote number larger than validators number") } votedAddrs := make([]bls.PublicKey, 0, validatorsBitSet.Count()) for index, val := range validators { if !validatorsBitSet.Test(uint(index)) { continue } voteAddr, err := bls.PublicKeyFromBytes(snap.Validators[val].VoteAddress[:]) if err != nil { return fmt.Errorf("BLS public key converts failed: %v", err) } votedAddrs = append(votedAddrs, voteAddr) } // The valid voted validators should be no less than 2/3 validators. if len(votedAddrs) < cmath.CeilDiv(len(snap.Validators)*2, 3) { return errors.New("invalid attestation, not enough validators voted") } // === Step 5: Signature verification === aggSig, err := bls.SignatureFromBytes(attestation.AggSignature[:]) if err != nil { return fmt.Errorf("BLS signature converts failed: %v", err) } if !aggSig.FastAggregateVerify(votedAddrs, attestation.Data.Hash()) { return errors.New("invalid attestation, signature verify failed") } return nil } // verifyHeader checks whether a header conforms to the consensus rules.The // caller may optionally pass in a batch of parents (ascending order) to avoid // looking those up from the database. This is useful for concurrently verifying // a batch of new headers. func (p *Parlia) verifyHeader(chain consensus.ChainHeaderReader, header *types.Header, parents []*types.Header) error { // Don't waste time checking blocks from the future if header.Time > uint64(time.Now().Unix()+time.Second.Milliseconds()/1000) { return consensus.ErrFutureBlock } if err := p.VerifyUnsealedHeader(chain, header, parents); err != nil { return err } // All basic checks passed, verify the seal and return return p.verifySeal(chain, header, parents) } // VerifyUnsealedHeader performs all header validity checks that do not require // a valid seal signature. It is used to validate a locally proposed block before // sealing: it runs the same structural, fork-rule, and cascading-field checks as // VerifyHeader but skips verifySeal (no signature yet) and verifyVoteAttestation // (vote attestation is embedded by the sealer and not present before sealing). func (p *Parlia) VerifyUnsealedHeader(chain consensus.ChainHeaderReader, header *types.Header, parents []*types.Header) error { // check extra data if len(header.Extra) < extraVanity { return errMissingVanity } if len(header.Extra) < extraVanity+extraSeal { return errMissingSignature } if header.Number == nil { return errUnknownBlock } number := header.Number.Uint64() epochLength, err := p.epochLength(chain, header, parents) if err != nil { return err } // Ensure that the extra-data contains a signer list on checkpoint, but none otherwise signersBytes := getValidatorBytesFromHeader(header, p.chainConfig, epochLength) isEpoch := number%epochLength == 0 if !isEpoch && len(signersBytes) != 0 && !bytes.HasPrefix(signersBytes, []byte("VFQ")) { return errExtraValidators } if isEpoch && len(signersBytes) == 0 { return errInvalidSpanValidators } lorentz := chain.Config().IsLorentz(header.Number, header.Time) if !lorentz { if header.MixDigest != (common.Hash{}) { return errInvalidMixDigest } } else { if header.MilliTimestamp()/1000 != header.Time { return fmt.Errorf("invalid MixDigest, have %#x, expected the last two bytes to represent milliseconds", header.MixDigest) } } // Ensure that the block doesn't contain any uncles which are meaningless in PoA if header.UncleHash != types.EmptyUncleHash { return errInvalidUncleHash } prague := chain.Config().IsPrague(header.Number, header.Time) if !prague { if header.ParentBeaconRoot != nil { return fmt.Errorf("invalid parentBeaconRoot, have %#x, expected nil", header.ParentBeaconRoot) } if header.RequestsHash != nil { return fmt.Errorf("invalid RequestsHash, have %#x, expected nil", header.RequestsHash) } } else { if header.ParentBeaconRoot == nil || *header.ParentBeaconRoot != (common.Hash{}) { return fmt.Errorf("invalid parentBeaconRoot, have %#x, expected zero hash", header.ParentBeaconRoot) } if header.RequestsHash == nil { return errors.New("header has nil RequestsHash after Prague") } } // All basic checks passed, verify cascading fields return p.verifyCascadingFields(chain, header, parents) } // verifyCascadingFields verifies all the header fields that are not standalone, // rather depend on a batch of previous headers. The caller may optionally pass // in a batch of parents (ascending order) to avoid looking those up from the // database. This is useful for concurrently verifying a batch of new headers. func (p *Parlia) verifyCascadingFields(chain consensus.ChainHeaderReader, header *types.Header, parents []*types.Header) error { // The genesis block is the always valid dead-end number := header.Number.Uint64() if number == 0 { return nil } parent, err := p.getParent(chain, header, parents) if err != nil { return err } // Now that the real parent is known, authoritatively validate the Snake8 // fork-activation data embedded in the header's extra-data. This rejects // forged embedded parent timestamps that would otherwise split consensus // between syncing and fully-synced nodes (BLK-3750). if err := p.verifySnake8Extra(header, parent); err != nil { return err } snap, err := p.snapshot(chain, number-1, header.ParentHash, parents, p.isSnake8Enabled(chain, header), header) if err != nil { return err } if _, ok := snap.Validators[header.Coinbase]; !ok { return errUnauthorizedValidator(header.Coinbase.String()) } if snap.SignRecently(header.Coinbase) { return errRecentlySigned } if header.Difficulty == nil { return errInvalidDifficulty } inturn := snap.inturn(header.Coinbase) if inturn && header.Difficulty.Cmp(diffInTurn) != 0 { return errWrongDifficulty } if !inturn && header.Difficulty.Cmp(diffNoTurn) != 0 { return errWrongDifficulty } if diff := new(big.Int).Sub(header.Number, parent.Number); diff.Cmp(big.NewInt(1)) != 0 { return consensus.ErrInvalidNumber } err = p.blockTimeVerifyForRamanujanFork(snap, header, parent) if err != nil { return err } // Verify the block's gas usage and (if applicable) verify the base fee. if !chain.Config().IsLondon(header.Number) { // Verify BaseFee not present before EIP-1559 fork. if header.BaseFee != nil { return fmt.Errorf("invalid baseFee before fork: have %d, expected 'nil'", header.BaseFee) } } else if err := eip1559.VerifyEIP1559Header(chain.Config(), parent, header); err != nil { // Verify the header's EIP-1559 attributes. return err } cancun := chain.Config().IsCancun(header.Number, header.Time) if !cancun { switch { case header.ExcessBlobGas != nil: return fmt.Errorf("invalid excessBlobGas: have %d, expected nil", header.ExcessBlobGas) case header.BlobGasUsed != nil: return fmt.Errorf("invalid blobGasUsed: have %d, expected nil", header.BlobGasUsed) case header.WithdrawalsHash != nil: return fmt.Errorf("invalid WithdrawalsHash, have %#x, expected nil", header.WithdrawalsHash) } } else { if !header.EmptyWithdrawalsHash() { return errors.New("header has wrong WithdrawalsHash") } if err := eip4844.VerifyEIP4844Header(chain.Config(), parent, header); err != nil { return err } } // Verify that the gas limit is <= 2^63-1 capacity := uint64(0x7fffffffffffffff) if header.GasLimit > capacity { return fmt.Errorf("invalid gasLimit: have %v, max %v", header.GasLimit, capacity) } // Verify that the gasUsed is <= gasLimit if header.GasUsed > header.GasLimit { return fmt.Errorf("invalid gasUsed: have %d, gasLimit %d", header.GasUsed, header.GasLimit) } // Verify that the gas limit remains within allowed bounds diff := int64(parent.GasLimit) - int64(header.GasLimit) if diff < 0 { diff *= -1 } gasLimitBoundDivisor := gasLimitBoundDivisorBeforeLorentz if p.chainConfig.IsLorentz(header.Number, header.Time) { gasLimitBoundDivisor = params.GasLimitBoundDivisor } limit := parent.GasLimit / gasLimitBoundDivisor if uint64(diff) >= limit || header.GasLimit < params.MinGasLimit { return fmt.Errorf("invalid gas limit: have %d, want %d += %d", header.GasLimit, parent.GasLimit, limit-1) } return nil } // snapshot retrieves the authorization snapshot at a given point in time. // !!! be careful // the block with `number` and `hash` is just the last element of `parents`, // unlike other interfaces such as verifyCascadingFields, `parents` are real parents func (p *Parlia) snapshot(chain consensus.ChainHeaderReader, number uint64, hash common.Hash, parents []*types.Header, isSnake8Fork bool, extraHeader *types.Header) (*Snapshot, error) { // Search for a snapshot in memory or on disk for checkpoints var ( headers []*types.Header snap *Snapshot ) for snap == nil { // If an in-memory snapshot was found, use that if s, ok := p.recentSnaps.Get(hash); ok { snap = s break } // If an on-disk checkpoint snapshot can be found, use that if number%checkpointInterval == 0 { if s, err := loadSnapshot(p.config, p.signatures, p.db, hash, p.ethAPI, isSnake8Fork); err == nil { log.Trace("Loaded snapshot from disk", "number", number, "hash", hash) snap = s break } } // If we're at the genesis, snapshot the initial state. Alternatively if we have // piled up more headers than allowed to be reorged (chain reinit from a freezer), // consider the checkpoint trusted and snapshot it. // Unable to retrieve the exact EpochLength here. // As known // defaultEpochLength = 200 && turnLength = 1 or 4 // lorentzEpochLength = 500 && turnLength = 8 // maxwellEpochLength = 1000 && turnLength = 16 // So just select block number like 1200, 2200, 3200, we can always get the right validators from `number - 200` offset := uint64(200) if number == 0 || (number%maxwellEpochLength == offset && (len(headers) > int(params.FullImmutabilityThreshold))) { var ( checkpoint *types.Header blockHash common.Hash blockInterval = defaultBlockInterval epochLength = defaultEpochLength ) if number == 0 { checkpoint = chain.GetHeaderByNumber(0) if checkpoint != nil { blockHash = checkpoint.Hash() } } else { checkpoint = chain.GetHeaderByNumber(number - offset) blockHeader := chain.GetHeaderByNumber(number) if blockHeader != nil { blockHash = blockHeader.Hash() if p.chainConfig.IsFermi(blockHeader.Number, blockHeader.Time) { blockInterval = fermiBlockInterval } else if p.chainConfig.IsMaxwell(blockHeader.Number, blockHeader.Time) { blockInterval = maxwellBlockInterval } else if p.chainConfig.IsLorentz(blockHeader.Number, blockHeader.Time) { blockInterval = lorentzBlockInterval } } if number > offset { // exclude `number == 200` blockBeforeCheckpoint := chain.GetHeaderByNumber(number - offset - 1) if blockBeforeCheckpoint != nil { if p.chainConfig.IsMaxwell(blockBeforeCheckpoint.Number, blockBeforeCheckpoint.Time) { epochLength = maxwellEpochLength } else if p.chainConfig.IsLorentz(blockBeforeCheckpoint.Number, blockBeforeCheckpoint.Time) { epochLength = lorentzEpochLength } } } } if checkpoint != nil && blockHash != (common.Hash{}) { // get validators from headers validators, voteAddrs, err := parseValidators(checkpoint, p.chainConfig, epochLength) if err != nil { return nil, err } // new snapshot snap = newSnapshot(p.config, p.signatures, number, blockHash, validators, voteAddrs, p.ethAPI, isSnake8Fork) // get turnLength from headers and use that for new turnLength turnLength, err := parseTurnLength(checkpoint, p.chainConfig, epochLength) if err != nil { return nil, err } if turnLength != nil { snap.TurnLength = *turnLength } snap.BlockInterval = blockInterval snap.EpochLength = epochLength // snap.Recents is currently empty, which affects the following: // a. The function SignRecently - This is acceptable since an empty snap.Recents results in a more lenient check. // b. The function blockTimeVerifyForRamanujanFork - This is also acceptable as it won't be invoked during `snap.apply`. // c. This may cause a mismatch in the slash systemtx, but the transaction list is not verified during `snap.apply`. // snap.Attestation is nil, but Snapshot.updateAttestation will handle it correctly. if err := snap.store(p.db); err != nil { return nil, err } log.Info("Stored checkpoint snapshot to disk", "number", number, "hash", blockHash) break } } // No snapshot for this header, gather the header and move backward var header *types.Header if len(parents) > 0 { // If we have explicit parents, pick from there (enforced) header = parents[len(parents)-1] if header.Hash() != hash || header.Number.Uint64() != number { return nil, consensus.ErrUnknownAncestor } parents = parents[:len(parents)-1] } else { // No explicit parents (or no more left), reach out to the database header = chain.GetHeader(hash, number) if header == nil { return nil, consensus.ErrUnknownAncestor } } headers = append(headers, header) number, hash = number-1, header.ParentHash } // check if snapshot is nil if snap == nil { return nil, fmt.Errorf("unknown error while retrieving snapshot at block number %v", number) } // Previous snapshot found, apply any pending headers on top of it for i := 0; i < len(headers)/2; i++ { headers[i], headers[len(headers)-1-i] = headers[len(headers)-1-i], headers[i] } snap, err := snap.apply(headers, chain, parents, p.chainConfig, isSnake8Fork) if err != nil { return nil, err } p.recentSnaps.Add(snap.Hash, snap) // If we've generated a new checkpoint snapshot, save to disk. Deliberately before // the Snake8 stamping below (COR-174) so nothing caller-dependent reaches disk. if snap.Number%checkpointInterval == 0 && len(headers) > 0 { if err = snap.store(p.db); err != nil { return nil, err } log.Trace("Stored snapshot to disk", "number", snap.Number, "hash", snap.Hash) } // FrequencyRLP is a per-caller view, so it goes on a copy the caller owns and never // on the shared cache entry (COR-174). Three reasons this matters: // - the verify path passes a not-yet-validated candidate header, which must not // be able to write into consensus state other code paths read; // - in-turn selection is derived from FrequencyRLP, so a cached value stamped by // an unrelated caller made the selected validator depend on call order; // - snapshot() is reached concurrently from header verification, the miner and // RPC handlers, so writing to the cached pointer was also a data race. // parlia_getSnapshot and friends (api.go) become read-only diagnostics by // construction: they receive a copy like everyone else. // // IsSnake8Fork and TurnLength are re-asserted here for the caller's view, but note // they are NOT caller-local — apply() and loadSnapshot() already installed them on // the cached entry, because the header loop depends on both. Do not "fix" that by // moving them here only; see the note in Snapshot.apply. // // Non-Snake8 callers have nothing to stamp and still get the shared entry, so the // standing rule for every caller is unchanged: treat the returned snapshot as // read-only unless you know it is a copy. if isSnake8Fork { snap = snap.copy() snap.IsSnake8Fork = true snap.TurnLength = snake8TurnLength if extraHeader != nil { if freq, err := parseValidatorFrequencies(extraHeader, p.chainConfig); err == nil { snap.FrequencyRLP = freq } } } var validators []string for v := range snap.Validators { validators = append(validators, v.Hex()) } log.Trace("loaded snapshot", "number", snap.Number, "hash", snap.Hash, "validators", strings.Join(validators, ","), "len", len(snap.Validators)) return snap, nil } // VerifyUncles implements consensus.Engine, always returning an error for any // uncles as this consensus mechanism doesn't permit uncles. func (p *Parlia) VerifyUncles(chain consensus.ChainReader, block *types.Block) error { if len(block.Uncles()) > 0 { return errors.New("uncles not allowed") } return nil } func (p *Parlia) VerifyRequests(header *types.Header, Requests [][]byte) error { return nil } // verifySeal checks whether the signature contained in the header satisfies the // consensus protocol requirements. The method accepts an optional list of parent // headers that aren't yet part of the local blockchain to generate the snapshots // from. func (p *Parlia) verifySeal(chain consensus.ChainHeaderReader, header *types.Header, parents []*types.Header) error { // Verifying the genesis block is not supported number := header.Number.Uint64() if number == 0 { return errUnknownBlock } // Verify vote attestation for fast finality. if err := p.verifyVoteAttestation(chain, header, parents); err != nil { log.Warn("Verify vote attestation failed", "error", err, "hash", header.Hash(), "number", header.Number, "parent", header.ParentHash, "coinbase", header.Coinbase, "extra", common.Bytes2Hex(header.Extra)) verifyVoteAttestationErrorCounter.Inc(1) if chain.Config().IsPlato(header.Number) { return err } } // Resolve the authorization key and check against validators signer, err := ecrecover(header, p.signatures, p.chainConfig.ChainID) if err != nil { return err } if signer != header.Coinbase { return errCoinBaseMisMatch } // check for double sign & add to cache key := proposalKey(*header) preHash, ok := p.recentHeaders.Get(key) if ok && preHash != header.Hash() { doubleSignCounter.Inc(1) log.Warn("DoubleSign detected", " block", header.Number, " miner", header.Coinbase, "hash1", preHash, "hash2", header.Hash()) } else { p.recentHeaders.Add(key, header.Hash()) } return nil } func (p *Parlia) prepareValidators(chain consensus.ChainHeaderReader, header *types.Header) error { epochLength, err := p.epochLength(chain, header, nil) if err != nil { return err } if header.Number.Uint64()%epochLength != 0 { return nil } newValidators, voteAddressMap, err := p.getCurrentValidators(header.ParentHash, new(big.Int).Sub(header.Number, big.NewInt(1))) if err != nil { return err } // sort validator by address sort.Sort(validatorsAscending(newValidators)) if !p.chainConfig.IsLuban(header.Number) { for _, validator := range newValidators { header.Extra = append(header.Extra, validator.Bytes()...) } } else { header.Extra = append(header.Extra, byte(len(newValidators))) if p.chainConfig.IsOnLuban(header.Number) { voteAddressMap = make(map[common.Address]*types.BLSPublicKey, len(newValidators)) var zeroBlsKey types.BLSPublicKey for _, validator := range newValidators { voteAddressMap[validator] = &zeroBlsKey } } for _, validator := range newValidators { header.Extra = append(header.Extra, validator.Bytes()...) header.Extra = append(header.Extra, voteAddressMap[validator].Bytes()...) } } return nil } func (p *Parlia) prepareTurnLength(chain consensus.ChainHeaderReader, header *types.Header) error { epochLength, err := p.epochLength(chain, header, nil) if err != nil { return err } if header.Number.Uint64()%epochLength != 0 || !p.chainConfig.IsBohr(header.Number, header.Time) { return nil } turnLength, err := p.getTurnLength(chain, header) if err != nil { return err } if turnLength != nil { header.Extra = append(header.Extra, *turnLength) } return nil } // assembleVoteAttestation collects votes and assembles the vote attestation into the block header. func (p *Parlia) assembleVoteAttestation(chain consensus.ChainHeaderReader, header *types.Header) error { // === Step 1: Preconditions === if !p.chainConfig.IsLuban(header.Number) || header.Number.Uint64() < 3 || p.VotePool == nil { return nil } // === Step 2: Find target header with quorum votes === parent := chain.GetHeaderByHash(header.ParentHash) if parent == nil { return errors.New("parent not found") } justifiedBlockNumber, justifiedBlockHash, err := p.GetJustifiedNumberAndHash(chain, []*types.Header{parent}) if err != nil { return errors.New("unexpected error when getting the highest justified number and hash") } var ( votes []*types.VoteEnvelope targetHeader = parent targetHeaderParentSnap *Snapshot ) for range p.GetAncestorGenerationDepth(header) { snap, err := p.snapshot(chain, targetHeader.Number.Uint64()-1, targetHeader.ParentHash, nil, p.isSnake8Enabled(chain, targetHeader), targetHeader) if err != nil { return err } votes = p.VotePool.FetchVotesByBlockHash(targetHeader.Hash(), justifiedBlockNumber) quorum := cmath.CeilDiv(len(snap.Validators)*2, 3) if len(votes) >= quorum { targetHeaderParentSnap = snap break } targetHeader = chain.GetHeaderByHash(targetHeader.ParentHash) if targetHeader == nil { return errors.New("parent not found") } if targetHeader.Number.Uint64() <= justifiedBlockNumber { break } } if targetHeaderParentSnap == nil { return nil } // === Step 3: Build vote attestation === attestation := &types.VoteAttestation{ Data: &types.VoteData{ SourceNumber: justifiedBlockNumber, SourceHash: justifiedBlockHash, TargetNumber: targetHeader.Number.Uint64(), TargetHash: targetHeader.Hash(), }, } // Validate vote data consistency for _, vote := range votes { if vote.Data.Hash() != attestation.Data.Hash() { return fmt.Errorf("vote check error, expected: %v, real: %v", attestation.Data, vote.Data) } } // Prepare aggregated vote signature voteAddrSet := make(map[types.BLSPublicKey]struct{}, len(votes)) signatures := make([][]byte, len(votes)) for i, vote := range votes { voteAddrSet[vote.VoteAddress] = struct{}{} signatures[i] = vote.Signature[:] } sigs, err := bls.MultipleSignaturesFromBytes(signatures) if err != nil { return err } copy(attestation.AggSignature[:], bls.AggregateSignatures(sigs).Marshal()) // Prepare vote address bitset. for _, valInfo := range targetHeaderParentSnap.Validators { if _, ok := voteAddrSet[valInfo.VoteAddress]; ok { attestation.VoteAddressSet |= 1 << (valInfo.Index - 1) // Index is offset by 1 } } bitsetCount := bitset.From([]uint64{uint64(attestation.VoteAddressSet)}).Count() if bitsetCount < uint(len(signatures)) { log.Warn(fmt.Sprintf("assembleVoteAttestation, check VoteAddress Set failed, expected:%d, real:%d", len(signatures), bitsetCount)) return errors.New("invalid attestation, check VoteAddress Set failed") } // === Step 4: Encode & insert into header extra === buf := new(bytes.Buffer) if err = rlp.Encode(buf, attestation); err != nil { return fmt.Errorf("attestation: failed to encode: %w", err) } extraSealStart := len(header.Extra) - extraSeal extraSealBytes := header.Extra[extraSealStart:] header.Extra = append(header.Extra[:extraSealStart], buf.Bytes()...) header.Extra = append(header.Extra, extraSealBytes...) return nil } // NextInTurnValidator return the next in-turn validator for header func (p *Parlia) NextInTurnValidator(chain consensus.ChainHeaderReader, header *types.Header) (common.Address, error) { snap, err := p.snapshot(chain, header.Number.Uint64(), header.Hash(), nil, p.isSnake8Enabled(chain, header), header) if err != nil { return common.Address{}, err } return snap.inturnValidator(), nil } // Prepare implements consensus.Engine, preparing all the consensus fields of the // header for running the transactions on top. func (p *Parlia) Prepare(chain consensus.ChainHeaderReader, header *types.Header) error { header.Coinbase = p.val return p.prepare(chain, header) } // prepare is shared by Prepare and PrepareForBidBlock; caller sets Coinbase/Number/ParentHash. func (p *Parlia) prepare(chain consensus.ChainHeaderReader, header *types.Header) error { header.Nonce = types.BlockNonce{} number := header.Number.Uint64() parent := chain.GetHeader(header.ParentHash, number-1) if parent == nil { return consensus.ErrUnknownAncestor } snap, err := p.snapshot(chain, number-1, header.ParentHash, nil, p.isSnake8Enabled(chain, header), nil) if err != nil { return err } // CONSENSUS-CRITICAL ORDERING (Chiliz Snake8): the difficulty stamped below // must be derived from the SAME frequency data that SetExtraData embeds in // this header — verifiers (and our own Seal) judge in-turn-ness from the // embedded data, so install the fresh FrequencyRLP on the snapshot BEFORE // computing difficulty and the Ramanujan backoff. Stamping difficulty from // the stale cached FrequencyRLP produces blocks that are invalid by // construction (diff=1 while the embedded data says in-turn) and forks the // chain — this regressed in the v1.7.6 prepare/SetExtraData split (COR-37). if p.isSnake8Enabled(chain, header) { if err := p.refreshFrequencyRLP(snap, number, parent); err != nil { // Post-Snake8Fix only: refuse to seal rather than emit a header // whose embedded bytes a verifier's recomputation would reject. return err } } // TODO: delete this log log.Trace("Prepare_start", "number", header.Number, "time", header.Time, "isSnake8", p.isSnake8Enabled(chain, header), "isSnake8Snap", snap.IsSnake8Fork, "inturnVal", snap.inturnValidator()) // Set the correct difficulty. header.Coinbase, not p.val: Prepare sets // Coinbase = p.val, but PrepareForBidBlock prepares for the in-turn // validator — the difficulty must describe the block's actual signer. header.Difficulty = calcDifficulty(snap, header.Coinbase) if header.Difficulty.Cmp(diffInTurn) != 0 && header.Number.Uint64() == 1 { return fmt.Errorf("not your turn for block producing") } // Ensure the extra data has all it's components if len(header.Extra) < extraVanity-nextForkHashSize { header.Extra = append(header.Extra, bytes.Repeat([]byte{0x00}, extraVanity-nextForkHashSize-len(header.Extra))...) } // Ensure the timestamp has the correct delay blockTime := p.blockTimeForRamanujanFork(snap, header, parent) header.Time = blockTime / 1000 // get seconds if p.chainConfig.IsLorentz(header.Number, header.Time) { header.SetMilliseconds(blockTime % 1000) } else { header.MixDigest = common.Hash{} } return p.SetExtraData(chain, header) } // refreshFrequencyRLP computes the fresh Snake8 validator-frequency data for // the child of snap (block `number`, whose parent header is `parent`) and installs // it on the snapshot. `snap` must be a caller-owned snapshot — Parlia.snapshot() // returns a copy under Snake8, so this write stays local and never reaches the // shared LRU entry (COR-174). // // Pre-Snake8Fix (historical behavior): stakes are read at the node's latest state // and every failure is soft — on a lookup error the stake degrades to zero, on a // calcFrequencyRLP error FrequencyRLP ends up nil and validator selection degrades // to round-robin. Because both prepare() (difficulty) and SetExtraData (embedded // data) go through this same helper, they degrade together and the produced header // stays self-consistent. Never returns an error on this path. // // Post-Snake8Fix (COR-173): the embedded bytes are consensus-verified, so they must // be reproducible by every node. Stake reads are pinned to the parent block's state // (by hash, so side chains resolve correctly) and a failed lookup is reported as // errSnake8StakeLookup — a NODE-LOCAL failure, not a statement about the block. // Sealing treats it as fatal (refuse to produce rather than embed bytes a verifier // would reject); verification fails open on it (see verifySnake8FrequencyData). The // one remaining soft failure is calcFrequencyRLP's "no eligible validators": it is a // deterministic function of the pinned stakes and snap.Recents, so producer and // verifier degrade to the same nil bytes (round-robin selection) and the header // still verifies. // // Pinned bytes are a pure function of the parent (stakes at the parent's state, // candidates from the parent's snapshot), so they are memoized in frequencyCache by // parent hash: prepare + SetExtraData + the Finalize verification all share one // stake sweep per parent. The unpinned pre-fork path must never be cached. func (p *Parlia) refreshFrequencyRLP(snap *Snapshot, number uint64, parent *types.Header) error { pinned := p.chainConfig.IsSnake8Fix(parent.Time) var state *rpc.BlockNumberOrHash if pinned { // Nil-guarded like stakeReader below: struct-literal engines have no cache. if p.frequencyCache != nil { if cached, ok := p.frequencyCache.Get(parent.Hash()); ok { snap.FrequencyRLP = bytes.Clone(cached) return nil } } parentState := rpc.BlockNumberOrHashWithHash(parent.Hash(), false) state = &parentState } // Fall back for engines assembled as struct literals (tests, tooling), which // bypass New() and leave the stakeReader seam nil. stakeReader := p.stakeReader if stakeReader == nil { stakeReader = p.getValidatorTotalDelegated } stakes := make(map[common.Address]*big.Int) for addr := range snap.Validators { totalDelegated, err := stakeReader(addr, number-1, state) if err != nil || totalDelegated == nil { if pinned { if err == nil { return fmt.Errorf("%w: validator %s at block %d: contract returned no stake", errSnake8StakeLookup, addr, number-1) } return fmt.Errorf("%w: validator %s at block %d: %v", errSnake8StakeLookup, addr, number-1, err) } log.Error("error when fetching total delegated amount", "validator", addr, "err", err) // Zero, not nil: calcFrequencyRLP dereferences every stake, so a // nil entry would crash the sealer on a transient RPC failure. // A zero stake excludes the validator from the frequency set, // which is the intended degradation. totalDelegated = big.NewInt(0) } stakes[addr] = totalDelegated } freqRlp, err := snap.calcFrequencyRLP(stakes) if err != nil { log.Error("error when calculating frequency rlp", "error", err, "block", number-1) } snap.FrequencyRLP = freqRlp if pinned && p.frequencyCache != nil { p.frequencyCache.Add(parent.Hash(), bytes.Clone(freqRlp)) } return nil } // SetExtraData rebuilds the validator-controlled extra-data section: // 28-byte vanity (caller-supplied, zero-padded) + 4-byte nextForkHash + // validators bytes (epoch only) + turnLength (epoch only, post-Bohr) + 65 reserved seal bytes. // Caller supplies the desired vanity in header.Extra; this function pads/truncates it. func (p *Parlia) SetExtraData(chain consensus.ChainHeaderReader, header *types.Header) error { // 32-byte vanity prefix: pad/truncate caller's vanity bytes, then append nextForkHash. if len(header.Extra) < extraVanity-nextForkHashSize { header.Extra = append(header.Extra, bytes.Repeat([]byte{0x00}, extraVanity-nextForkHashSize-len(header.Extra))...) } header.Extra = header.Extra[:extraVanity-nextForkHashSize] nextForkHash := forkid.NextForkHash(p.chainConfig, p.genesisHash, chain.GenesisHeader().Time, header.Number.Uint64(), header.Time) header.Extra = append(header.Extra, nextForkHash[:]...) if err := p.prepareValidators(chain, header); err != nil { return err } if err := p.prepareTurnLength(chain, header); err != nil { return err } // Add RLP-encoded validator+frequency data (Chiliz Snake8). SetExtraData is // also called standalone (e.g. from the miner MEV path), so recompute the // parent and the frequency snapshot here rather than relying on prepare's // locals. refreshFrequencyRLP is deterministic for a given parent, so this // recomputation yields exactly the bytes prepare() derived the difficulty // from — the header stays coherent (see the ordering note in prepare). if p.isSnake8Enabled(chain, header) { number := header.Number.Uint64() parent := chain.GetHeader(header.ParentHash, number-1) if parent == nil { return consensus.ErrUnknownAncestor } snap, err := p.snapshot(chain, number-1, header.ParentHash, nil, true, nil) if err != nil { return err } if err := p.refreshFrequencyRLP(snap, number, parent); err != nil { return err } ts := make([]byte, 8) binary.LittleEndian.PutUint64(ts, parent.Time) log.Trace("Prepare", "append snake8 data", parent.Time, "number", header.Number.Uint64(), "len(validatorFrequencyDataPrefix)", len(validatorFrequencyDataPrefix), "len(ts)", len(ts), "len(snap.FrequencyRLP)", len(snap.FrequencyRLP)) header.Extra = append(header.Extra, validatorFrequencyDataPrefix...) header.Extra = append(header.Extra, ts...) header.Extra = append(header.Extra, snap.FrequencyRLP...) } // add extra seal space header.Extra = append(header.Extra, make([]byte, extraSeal)...) return nil } func (p *Parlia) verifyValidators(chain consensus.ChainHeaderReader, header *types.Header) error { epochLength, err := p.epochLength(chain, header, nil) if err != nil { return err } if header.Number.Uint64()%epochLength != 0 { return nil } newValidators, voteAddressMap, err := p.getCurrentValidators(header.ParentHash, new(big.Int).Sub(header.Number, big.NewInt(1))) if err != nil { return err } // sort validator by address sort.Sort(validatorsAscending(newValidators)) var validatorsBytes []byte validatorsNumber := len(newValidators) if !p.chainConfig.IsLuban(header.Number) { validatorsBytes = make([]byte, validatorsNumber*validatorBytesLengthBeforeLuban) for i, validator := range newValidators { copy(validatorsBytes[i*validatorBytesLengthBeforeLuban:], validator.Bytes()) } } else { if uint8(validatorsNumber) != header.Extra[extraVanity] { return errMismatchingEpochValidators } validatorsBytes = make([]byte, validatorsNumber*validatorBytesLength) if p.chainConfig.IsOnLuban(header.Number) { voteAddressMap = make(map[common.Address]*types.BLSPublicKey, len(newValidators)) var zeroBlsKey types.BLSPublicKey for _, validator := range newValidators { voteAddressMap[validator] = &zeroBlsKey } } for i, validator := range newValidators { copy(validatorsBytes[i*validatorBytesLength:], validator.Bytes()) copy(validatorsBytes[i*validatorBytesLength+common.AddressLength:], voteAddressMap[validator].Bytes()) } } if !bytes.Equal(getValidatorBytesFromHeader(header, p.chainConfig, epochLength), validatorsBytes) { return errMismatchingEpochValidators } return nil } func (p *Parlia) verifyTurnLength(chain consensus.ChainHeaderReader, header *types.Header) error { epochLength, err := p.epochLength(chain, header, nil) if err != nil { return err } if header.Number.Uint64()%epochLength != 0 || !p.chainConfig.IsBohr(header.Number, header.Time) { return nil } turnLengthFromHeader, err := parseTurnLength(header, p.chainConfig, epochLength) if err != nil { return err } if turnLengthFromHeader != nil { turnLength, err := p.getTurnLength(chain, header) if err != nil { return err } if turnLength != nil && *turnLength == *turnLengthFromHeader { log.Debug("verifyTurnLength", "turnLength", *turnLength) return nil } } return errMismatchingEpochTurnLength } func (p *Parlia) distributeFinalityReward(chain consensus.ChainHeaderReader, state vm.StateDB, header *types.Header, cx core.ChainContext, txs *[]*types.Transaction, receipts *[]*types.Receipt, systemTxs *[]*types.Transaction, usedGas *uint64, mode systemTxMode, tracer *tracing.Hooks) error { currentHeight := header.Number.Uint64() if currentHeight%finalityRewardInterval != 0 { return nil } head := header accumulatedWeights := make(map[common.Address]uint64) for height := currentHeight - 1; height+finalityRewardInterval >= currentHeight && height >= 1; height-- { head = chain.GetHeaderByHash(head.ParentHash) if head == nil { return fmt.Errorf("header is nil at height %d", height) } epochLength, err := p.epochLength(chain, head, nil) if err != nil { return err } voteAttestation, err := getVoteAttestationFromHeader(head, chain.Config(), epochLength) if err != nil { return err } if voteAttestation == nil { continue } justifiedBlock := chain.GetHeaderByHash(voteAttestation.Data.TargetHash) if justifiedBlock == nil { log.Warn("justifiedBlock is nil at height %d", voteAttestation.Data.TargetNumber) continue } parent := chain.GetHeaderByHash(justifiedBlock.ParentHash) if parent == nil { return consensus.ErrUnknownAncestor } snap, err := p.snapshot(chain, justifiedBlock.Number.Uint64()-1, justifiedBlock.ParentHash, nil, p.isSnake8Enabled(chain, head), head) if err != nil { return err } validators := snap.validators() validatorsBitSet := bitset.From([]uint64{uint64(voteAttestation.VoteAddressSet)}) if validatorsBitSet.Count() > uint(len(validators)) { log.Error("invalid attestation, vote number larger than validators number") continue } validVoteCount := 0 for index, val := range validators { if validatorsBitSet.Test(uint(index)) { accumulatedWeights[val] += 1 validVoteCount += 1 } } quorum := cmath.CeilDiv(len(snap.Validators)*2, 3) if validVoteCount > quorum { accumulatedWeights[head.Coinbase] += uint64((validVoteCount - quorum) * collectAdditionalVotesRewardRatio / 100) } } validators := make([]common.Address, 0, len(accumulatedWeights)) weights := make([]*big.Int, 0, len(accumulatedWeights)) for val := range accumulatedWeights { validators = append(validators, val) } sort.Sort(validatorsAscending(validators)) for _, val := range validators { weights = append(weights, big.NewInt(int64(accumulatedWeights[val]))) } // generate system transaction method := "distributeFinalityReward" data, err := p.validatorSetABI.Pack(method, validators, weights) if err != nil { log.Error("Unable to pack tx for distributeFinalityReward", "error", err) return err } msg := p.getSystemMessage(header.Coinbase, common.HexToAddress(systemcontracts.ValidatorContract), data, common.Big0) return p.applyTransaction(msg, state, header, cx, txs, receipts, systemTxs, usedGas, mode, tracer) } func (p *Parlia) EstimateGasReservedForSystemTxs(chain consensus.ChainHeaderReader, header *types.Header) uint64 { parent := chain.GetHeaderByHash(header.ParentHash) if parent != nil { // Mainnet and Chapel have both passed Feynman. Now, simplify the logic before and during the Feynman hard fork. if p.chainConfig.IsFeynman(header.Number, header.Time) && !p.chainConfig.IsOnFeynman(header.Number, parent.Time, header.Time) { // const ( // the following values represent the maximum values found in the most recent blocks on the mainnet // depositTxGas = uint64(60_000) // slashTxGas = uint64(140_000) // finalityRewardTxGas = uint64(350_000) // updateValidatorTxGas = uint64(12_160_000) // ) // suggestReservedGas := depositTxGas // if header.Difficulty.Cmp(diffInTurn) != 0 { // snap, err := p.snapshot(chain, header.Number.Uint64()-1, header.ParentHash, nil) // if err != nil || !snap.SignRecently(snap.inturnValidator()) { // suggestReservedGas += slashTxGas // } // } // if header.Number.Uint64()%p.config.Epoch == 0 { // suggestReservedGas += finalityRewardTxGas // } // if isBreatheBlock(parent.Time, header.Time) { // suggestReservedGas += updateValidatorTxGas // } // return suggestReservedGas * 150 / 100 if !isBreatheBlock(parent.Time, header.Time) { // params.SystemTxsGasSoftLimit > (depositTxGas+slashTxGas+finalityRewardTxGas)*150/100 return params.SystemTxsGasSoftLimit } } } // params.SystemTxsGasHardLimit > (depositTxGas+slashTxGas+finalityRewardTxGas+updateValidatorTxGas)*150/100 return params.SystemTxsGasHardLimit } // Finalize implements consensus.Engine, ensuring no uncles are set, nor block // rewards given. func (p *Parlia) Finalize(chain consensus.ChainHeaderReader, header *types.Header, state vm.StateDB, txs *[]*types.Transaction, uncles []*types.Header, _ []*types.Withdrawal, receipts *[]*types.Receipt, systemTxs *[]*types.Transaction, usedGas *uint64, tracer *tracing.Hooks) error { // warn if not in majority fork p.detectNewVersionWithFork(chain, header, state) // If the block is an epoch end block, verify the validator list // The verification can only be done when the state is ready, it can't be done in VerifyHeader. if err := p.verifyValidators(chain, header); err != nil { return err } if err := p.verifyTurnLength(chain, header); err != nil { return err } cx := chainContext{ChainHeaderReader: chain, parlia: p} parent := chain.GetHeaderByHash(header.ParentHash) if parent == nil { return consensus.ErrUnknownAncestor } // Snake8Fix (COR-173): with state now available, verify the embedded // validator-frequency data against its recomputation from the parent // block's state. Like verifyValidators above, this cannot run in // VerifyHeader because it needs a contract call against parent state. if err := p.verifySnake8FrequencyData(chain, header, parent); err != nil { return err } systemcontracts.TryUpdateBuildInSystemContract(p.chainConfig, header.Number, parent.Time, header.Time, state, false) if p.chainConfig.IsOnFeynman(header.Number, parent.Time, header.Time) { err := p.initializeFeynmanContract(state, header, cx, txs, receipts, systemTxs, usedGas, systemTxImporting, tracer) if err != nil { return fmt.Errorf("init feynman contract failed: %v", err) } } // No block rewards in PoA, so the state remains as is and uncles are dropped if header.Number.Cmp(common.Big1) == 0 { err := p.initContract(state, header, cx, txs, receipts, systemTxs, usedGas, systemTxImporting, tracer) if err != nil { log.Error("init contract failed", "error", err) return err } } if header.Difficulty.Cmp(diffInTurn) != 0 { snap, err := p.snapshot(chain, header.Number.Uint64()-1, header.ParentHash, nil, p.isSnake8Enabled(chain, header), header) if err != nil { return err } spoiledVal := snap.inturnValidator() signedRecently := false if p.chainConfig.IsPlato(header.Number) { signedRecently = snap.SignRecently(spoiledVal) } else { for _, recent := range snap.Recents { if recent == spoiledVal { signedRecently = true break } } } if !signedRecently { log.Trace("slash validator", "block hash", header.Hash(), "address", spoiledVal) err = p.slash(spoiledVal, state, header, cx, txs, receipts, systemTxs, usedGas, systemTxImporting, tracer) if err != nil { log.Error("slash validator failed", "block hash", header.Hash(), "address", spoiledVal, "err", err) } } } val := header.Coinbase PenalizeForDelayMining, err := p.isIntentionalDelayMining(chain, header) if err != nil { log.Debug("unexpected error happened when detecting intentional delay mining", "err", err) } if PenalizeForDelayMining { intentionalDelayMiningCounter.Inc(1) log.Warn("intentional delay mining detected", "validator", val, "number", header.Number, "hash", header.Hash()) } err = p.distributeIncoming(val, state, header, cx, txs, receipts, systemTxs, usedGas, systemTxImporting, tracer) if err != nil { return err } if p.chainConfig.IsPlato(header.Number) { if err := p.distributeFinalityReward(chain, state, header, cx, txs, receipts, systemTxs, usedGas, systemTxImporting, tracer); err != nil { return err } } // update validators every day if p.chainConfig.IsFeynman(header.Number, header.Time) && isBreatheBlock(parent.Time, header.Time) { // we should avoid update validators in the Feynman upgrade block if !p.chainConfig.IsOnFeynman(header.Number, parent.Time, header.Time) { if err := p.updateValidatorSetV2(state, header, cx, txs, receipts, systemTxs, usedGas, systemTxImporting, tracer); err != nil { return err } } } if len(*systemTxs) > 0 { return errors.New("the length of systemTxs do not match") } return nil } type systemTxMode uint8 const ( systemTxImporting systemTxMode = iota systemTxMining systemTxPacking ) func (p *Parlia) FinalizeAndAssemble(chain consensus.ChainHeaderReader, header *types.Header, state *state.StateDB, body *types.Body, receipts []*types.Receipt, tracer *tracing.Hooks) (*types.Block, []*types.Receipt, error) { return p.finalizeAndAssemble(chain, header, state, body, receipts, tracer, systemTxMining) } func (p *Parlia) finalizeAndAssemble(chain consensus.ChainHeaderReader, header *types.Header, state *state.StateDB, body *types.Body, receipts []*types.Receipt, tracer *tracing.Hooks, mode systemTxMode) (*types.Block, []*types.Receipt, error) { // No block rewards in PoA, so the state remains as is and uncles are dropped cx := chainContext{ChainHeaderReader: chain, parlia: p} if body.Transactions == nil { body.Transactions = make([]*types.Transaction, 0) } if receipts == nil { receipts = make([]*types.Receipt, 0) } parent := chain.GetHeaderByHash(header.ParentHash) if parent == nil { return nil, nil, errors.New("parent not found") } systemcontracts.TryUpdateBuildInSystemContract(p.chainConfig, header.Number, parent.Time, header.Time, state, false) if p.chainConfig.IsOnFeynman(header.Number, parent.Time, header.Time) { err := p.initializeFeynmanContract(state, header, cx, &body.Transactions, &receipts, nil, &header.GasUsed, mode, tracer) if err != nil { return nil, nil, fmt.Errorf("init feynman contract failed: %v", err) } } if header.Number.Cmp(common.Big1) == 0 { err := p.initContract(state, header, cx, &body.Transactions, &receipts, nil, &header.GasUsed, mode, tracer) if err != nil { log.Error("init contract failed", "error", err) return nil, nil, err } } if header.Difficulty.Cmp(diffInTurn) != 0 { number := header.Number.Uint64() snap, err := p.snapshot(chain, number-1, header.ParentHash, nil, p.isSnake8Enabled(chain, header), header) if err != nil { return nil, nil, err } spoiledVal := snap.inturnValidator() signedRecently := false if p.chainConfig.IsPlato(header.Number) { signedRecently = snap.SignRecently(spoiledVal) } else { for _, recent := range snap.Recents { if recent == spoiledVal { signedRecently = true break } } } if !signedRecently { err = p.slash(spoiledVal, state, header, cx, &body.Transactions, &receipts, nil, &header.GasUsed, mode, tracer) if err != nil { log.Error("slash validator failed", "block hash", header.Hash(), "address", spoiledVal) } } } err := p.distributeIncoming(header.Coinbase, state, header, cx, &body.Transactions, &receipts, nil, &header.GasUsed, mode, tracer) if err != nil { return nil, nil, err } if p.chainConfig.IsPlato(header.Number) { if err := p.distributeFinalityReward(chain, state, header, cx, &body.Transactions, &receipts, nil, &header.GasUsed, mode, tracer); err != nil { return nil, nil, err } } // update validators every day if p.chainConfig.IsFeynman(header.Number, header.Time) && isBreatheBlock(parent.Time, header.Time) { // we should avoid update validators in the Feynman upgrade block if !p.chainConfig.IsOnFeynman(header.Number, parent.Time, header.Time) { if err := p.updateValidatorSetV2(state, header, cx, &body.Transactions, &receipts, nil, &header.GasUsed, mode, tracer); err != nil { return nil, nil, err } } } // should not happen. Once happen, stop the node is better than broadcast the block if header.GasLimit < header.GasUsed { return nil, nil, errors.New("gas consumption of system txs exceed the gas limit") } header.UncleHash = types.EmptyUncleHash var blk *types.Block var rootHash common.Hash wg := sync.WaitGroup{} wg.Add(2) go func() { rootHash = state.IntermediateRoot(chain.Config().IsEIP158(header.Number)) wg.Done() }() go func() { blk = types.NewBlock(header, body, receipts, trie.NewStackTrie(nil)) wg.Done() }() wg.Wait() blk.SetRoot(rootHash) // Assemble and return the final block for sealing return blk, receipts, nil } func (p *Parlia) IsActiveValidatorAt(chain consensus.ChainHeaderReader, header *types.Header, checkVoteKeyFn func(bLSPublicKey *types.BLSPublicKey) bool) bool { number := header.Number.Uint64() snap, err := p.snapshot(chain, number-1, header.ParentHash, nil, p.isSnake8Enabled(chain, header), header) if err != nil { log.Error("failed to get the snapshot from consensus", "error", err) return false } validators := snap.Validators validatorInfo, ok := validators[p.val] return ok && (checkVoteKeyFn == nil || (validatorInfo != nil && checkVoteKeyFn(&validatorInfo.VoteAddress))) } // VerifyVote will verify: 1. If the vote comes from valid validators 2. If the vote's sourceNumber and sourceHash are correct func (p *Parlia) VerifyVote(chain consensus.ChainHeaderReader, vote *types.VoteEnvelope) error { targetNumber := vote.Data.TargetNumber targetHash := vote.Data.TargetHash header := chain.GetVerifiedBlockByHash(targetHash) if header == nil { log.Warn("BlockHeader at current voteBlockNumber is nil", "targetNumber", targetNumber, "targetHash", targetHash) return errors.New("BlockHeader at current voteBlockNumber is nil") } if header.Number.Uint64() != targetNumber { log.Warn("unexpected target number", "expect", header.Number.Uint64(), "real", targetNumber) return errors.New("target number mismatch") } justifiedBlockNumber, justifiedBlockHash, err := p.GetJustifiedNumberAndHash(chain, []*types.Header{header}) if err != nil { log.Error("failed to get the highest justified number and hash", "headerNumber", header.Number, "headerHash", header.Hash()) return errors.New("unexpected error when getting the highest justified number and hash") } if vote.Data.SourceNumber != justifiedBlockNumber || vote.Data.SourceHash != justifiedBlockHash { return errors.New("vote source block mismatch") } number := header.Number.Uint64() snap, err := p.snapshot(chain, number-1, header.ParentHash, nil, p.isSnake8Enabled(chain, header), header) if err != nil { log.Error("failed to get the snapshot from consensus", "error", err) return errors.New("failed to get the snapshot from consensus") } validators := snap.Validators voteAddress := vote.VoteAddress for addr, validator := range validators { if validator.VoteAddress == voteAddress { if addr == p.val { validVotesfromSelfCounter.Inc(1) } metrics.GetOrRegisterCounter(fmt.Sprintf("parlia/VerifyVote/%s", addr.String()), nil).Inc(1) return nil } } return errors.New("vote verification failed") } // Authorize injects a private key into the consensus engine to mint new blocks // with. func (p *Parlia) Authorize(val common.Address, signFn SignerFn, signTxFn SignerTxFn) { p.lock.Lock() defer p.lock.Unlock() p.val = val p.signFn = signFn p.signTxFn = signTxFn } // IsLastBlockInTurn reports whether header is the last block in the current validator's turn. func (p *Parlia) IsLastBlockInTurn(chain consensus.ChainReader, header *types.Header) bool { snap, err := p.snapshot(chain, header.Number.Uint64()-1, header.ParentHash, nil, p.isSnake8Enabled(chain, header), header) if err != nil { return false } return snap.lastBlockInOneTurn(header.Number.Uint64()) } // Argument leftOver is the time reserved for block finalize(calculate root, distribute income...) func (p *Parlia) Delay(chain consensus.ChainReader, header *types.Header, leftOver *time.Duration) *time.Duration { snap, err := p.snapshot(chain, header.Number.Uint64()-1, header.ParentHash, nil, p.isSnake8Enabled(chain, header), header) if err != nil { return nil } delay := p.delayForRamanujanFork(snap, header) timeForMining := time.Duration(snap.BlockInterval) * time.Millisecond if delay > timeForMining { delay = timeForMining } if *leftOver >= time.Duration(snap.BlockInterval)*time.Millisecond { // ignore invalid leftOver log.Error("Delay invalid argument", "leftOver", leftOver.String(), "Period", snap.BlockInterval) } else if *leftOver >= delay { delay = time.Duration(0) return &delay } else { delay = delay - *leftOver } return &delay } // Seal implements consensus.Engine, attempting to create a sealed block using // the local signing credentials. func (p *Parlia) Seal(chain consensus.ChainHeaderReader, block *types.Block, results chan<- *types.Block, stop <-chan struct{}) error { header := block.Header() // Sealing the genesis block is not supported number := header.Number.Uint64() if number == 0 { return errUnknownBlock } // Don't hold the val fields for the entire sealing procedure p.lock.RLock() val, signFn := p.val, p.signFn p.lock.RUnlock() isSnake8 := p.isSnake8Enabled(chain, header) snap, err := p.snapshot(chain, number-1, header.ParentHash, nil, isSnake8, header) if err != nil { return err } // Bail out if we're unauthorized to sign a block if _, authorized := snap.Validators[val]; !authorized { return errUnauthorizedValidator(val.String()) } // If we're amongst the recent signers, wait for the next block if snap.SignRecently(val) { log.Info("Signed recently, must wait for others") return nil } // Sweet, the protocol permits us to sign the block, wait for our time delay := p.delayForRamanujanFork(snap, header) log.Info("Sealing block with", "number", number, "delay", delay, "headerDifficulty", header.Difficulty, "val", val.Hex(), "isSnake8", isSnake8, "isSnake8Snap", snap.IsSnake8Fork, "inturnVal", snap.inturnValidator()) // Wait until sealing is terminated or delay timeout. log.Info("Waiting for slot to sign and propagate", "delay", common.PrettyDuration(delay)) go func() { select { case <-stop: return case <-time.After(delay): } err := p.assembleVoteAttestation(chain, header) if err != nil { /* If the vote attestation can't be assembled successfully, the blockchain won't get fast finalized, but it can be tolerated, so just report this error here. */ log.Debug("Assemble vote attestation failed when sealing", "err", err) } // Sign all the things! sig, err := signFn(accounts.Account{Address: val}, accounts.MimetypeParlia, ParliaRLP(header, p.chainConfig.ChainID)) if err != nil { log.Error("Sign for the block header failed when sealing", "err", err) return } copy(header.Extra[len(header.Extra)-extraSeal:], sig) if p.shouldWaitForCurrentBlockProcess(chain, header, snap) { highestVerifiedHeader := chain.GetHighestVerifiedHeader() // including time for writing and committing blocks waitProcessEstimate := math.Ceil(float64(highestVerifiedHeader.GasUsed) / float64(100_000_000)) log.Info("Waiting for received in turn block to process", "waitProcessEstimate(Seconds)", waitProcessEstimate) select { case <-stop: log.Info("Received block process finished, abort block seal") return case <-time.After(time.Duration(waitProcessEstimate) * time.Second): if chain.CurrentHeader().Number.Uint64() >= header.Number.Uint64() { log.Info("Process backoff time exhausted, and current header has updated to abort this seal") return } log.Info("Process backoff time exhausted, start to seal block") } } select { case results <- block.WithSeal(header): default: log.Warn("Sealing result is not read by miner", "sealhash", types.SealHash(header, p.chainConfig.ChainID)) } }() return nil } func (p *Parlia) shouldWaitForCurrentBlockProcess(chain consensus.ChainHeaderReader, header *types.Header, snap *Snapshot) bool { if header.Difficulty.Cmp(diffInTurn) == 0 { return false } highestVerifiedHeader := chain.GetHighestVerifiedHeader() if highestVerifiedHeader == nil { log.Warn("Shouldn't wait for block process, because there is no highest verified header") return false } if header.ParentHash == highestVerifiedHeader.ParentHash { return true } return false } func (p *Parlia) EnoughDistance(chain consensus.ChainReader, header *types.Header) bool { snap, err := p.snapshot(chain, header.Number.Uint64()-1, header.ParentHash, nil, p.isSnake8Enabled(chain, header), header) if err != nil { return true } return snap.enoughDistance(p.val, header) } func (p *Parlia) IsLocalBlock(header *types.Header) bool { return p.val == header.Coinbase } func (p *Parlia) SignRecently(chain consensus.ChainReader, parent *types.Header) (bool, error) { snap, err := p.snapshot(chain, parent.Number.Uint64(), parent.Hash(), nil, p.isSnake8Enabled(chain, parent), parent) if err != nil { return true, err } // Bail out if we're unauthorized to sign a block if _, authorized := snap.Validators[p.val]; !authorized { return true, errUnauthorizedValidator(p.val.String()) } return snap.SignRecently(p.val), nil } // CalcDifficulty is the difficulty adjustment algorithm. It returns the difficulty // that a new block should have based on the previous blocks in the chain and the // current signer. // // Snake8 in-turn rule (COR-173): block N's in-turn-ness — and therefore its // consensus-verified difficulty — is judged from the frequency bytes embedded in // header N itself (verifyCascadingFields stamps the candidate header's bytes onto // the snapshot). Post-Snake8Fix those bytes are in turn forced to equal the // recomputation from state at N-1 (verifySnake8FrequencyData), which is what the // producer embeds via prepare()/SetExtraData. This function is only a scheduling // hint for the miner: it approximates block N+1's difficulty using the bytes // embedded in parent header N (state N-1 derived) rather than the state-N-derived // bytes the actual block will carry. That approximation is not consensus-relevant // — Prepare recomputes the authoritative difficulty before sealing. func (p *Parlia) CalcDifficulty(chain consensus.ChainHeaderReader, time uint64, parent *types.Header) *big.Int { snap, err := p.snapshot(chain, parent.Number.Uint64(), parent.Hash(), nil, p.isSnake8Enabled(chain, parent), parent) if err != nil { return nil } return calcDifficulty(snap, p.val) } // CalcDifficulty is the difficulty adjustment algorithm. It returns the difficulty // that a new block should have based on the previous blocks in the chain and the // current signer. func calcDifficulty(snap *Snapshot, signer common.Address) *big.Int { if snap.inturn(signer) { return new(big.Int).Set(diffInTurn) } return new(big.Int).Set(diffNoTurn) } func encodeSigHeaderWithoutVoteAttestation(w io.Writer, header *types.Header, chainId *big.Int) { err := rlp.Encode(w, []interface{}{ chainId, header.ParentHash, header.UncleHash, header.Coinbase, header.Root, header.TxHash, header.ReceiptHash, header.Bloom, header.Difficulty, header.Number, header.GasLimit, header.GasUsed, header.Time, header.Extra[:extraVanity], // this will panic if extra is too short, should check before calling encodeSigHeaderWithoutVoteAttestation header.MixDigest, header.Nonce, }) if err != nil { panic("can't encode: " + err.Error()) } } // SealHash returns the hash of a block without vote attestation prior to it being sealed. // So it's not the real hash of a block, just used as unique id to distinguish task func (p *Parlia) SealHash(header *types.Header) (hash common.Hash) { hasher := sha3.NewLegacyKeccak256() encodeSigHeaderWithoutVoteAttestation(hasher, header, p.chainConfig.ChainID) hasher.Sum(hash[:0]) return hash } // APIs implements consensus.Engine, returning the user facing RPC API to query snapshot. func (p *Parlia) APIs(chain consensus.ChainHeaderReader) []rpc.API { return []rpc.API{{ Namespace: "parlia", Version: "1.0", Service: &API{chain: chain, parlia: p}, Public: false, }} } // Close implements consensus.Engine. It's a noop for parlia as there are no background threads. func (p *Parlia) Close() error { return nil } // ========================== interaction with contract/account ========= // Returns the inflation % for years passed // years are estimated based on the seconds passed since the fork block's timestamp // If year > 13: pct = 1.88, Otherwise pct = 9.24e^(-0.250x) + 1.60 // Note: result is rounded to 1e18 decimal places and returned as int (percent * 1e18) func getInflationPct(secondsPassed uint64) *big.Int { // convert seconds to years year := big.NewFloat(0) year.Mul(big.NewFloat(0).SetUint64(secondsPassed), big.NewFloat(1.0/31536000)) year.Add(year, big.NewFloat(1)) yearF, _ := year.Float64() log.Trace("inflation year", "year", yearF) var inflationPct float64 if year.Cmp(big.NewFloat(13)) > 0 { inflationPct = 1.88 } else { initDecayMag := 9.24 decayRate := -0.25 offset := 1.6 inflationPct = initDecayMag*math.Pow(math.E, decayRate*yearF) + offset } inflationPct = math.Round(inflationPct*1e18) / 1e18 return big.NewInt(int64(inflationPct * 1e18)) } // Returns amount of chz & inflation pct for specific block (part of dragon8) func getNewSupplyForBlock(forkTs uint64, currentTs uint64, lastSupply *big.Int) (*big.Int, *big.Int) { // Calculate the new supply secondsPassed := currentTs - forkTs newIntroducedSupply := big.NewInt(0) inflationPct := getInflationPct(secondsPassed) newIntroducedSupply.Mul(lastSupply, inflationPct) newIntroducedSupply.Div(newIntroducedSupply, big.NewInt(1e18)) newIntroducedSupply.Div(newIntroducedSupply, big.NewInt(100)) // inflPct is percent*1e18 // Calculate amount for 1 block blockAmount := big.NewInt(0) blockAmount.Div(newIntroducedSupply, big.NewInt(10512000)) // 10512000 blocks = ~1y log.Trace("inflation details", "secondsSinceFork", secondsPassed, "newIntroducedSupply", newIntroducedSupply, "blockAmount", blockAmount) return blockAmount, inflationPct } // Returns inflation %, supply, amount for the block (part of dragon8Fix) func getNewSupplyForBlockDragon8Fix(forkTime uint64, currentTime uint64) (*big.Int, *big.Int, *big.Int) { var ( // inflation %, supply, amount per block inflationData = [][]*big.Int{ {big.NewInt(87961192355797800), cmath.MustParseBig256("8888888888000000000000000000"), cmath.MustParseBig256("74379496319128800000")}, // year 0 {big.NewInt(72043432957447300), cmath.MustParseBig256("9670766153000000000000000000"), cmath.MustParseBig256("66278081527102500000")}, // year 1 // one time inflation of 148600000 CHZ during year 1 + og schedule // new supply after year1 = 10367481346 + 148600000 = 10516081346 // year 3 to year 8 have changed {big.NewInt(55304937824698300), cmath.MustParseBig256("10516081346000000000000000000"), cmath.MustParseBig256("55326410292998500000")}, // year 2 {big.NewInt(46543801590000000), cmath.MustParseBig256("11097672571000000000000000000"), cmath.MustParseBig256("49136973934551000000")}, // year 3 {big.NewInt(39673708820000000), cmath.MustParseBig256("11614200441000000000000000000"), cmath.MustParseBig256("43833562214611900000")}, // year 4 {big.NewInt(39673708820000000), cmath.MustParseBig256("12074978847000000000000000000"), cmath.MustParseBig256("39395232686453600000")}, // year 5 {big.NewInt(34295934680000000), cmath.MustParseBig256("12489101533000000000000000000"), cmath.MustParseBig256("35751611491628600000")}, // year 6 {big.NewInt(30091911660000000), cmath.MustParseBig256("12864922473000000000000000000"), cmath.MustParseBig256("34885308219178100000")}, // year 7 // {big.NewInt(25738888349516300), cmath.MustParseBig256("13231636833000000000000000000"), cmath.MustParseBig256("32397985455982600000")}, // year 8 {big.NewInt(23584653872848300), cmath.MustParseBig256("13572204456000000000000000000"), cmath.MustParseBig256("30450508407284500000")}, // year 9 {big.NewInt(21906934375499800), cmath.MustParseBig256("13892300200000000000000000000"), cmath.MustParseBig256("28951456317173600000")}, // year 10 {big.NewInt(20600325117190600), cmath.MustParseBig256("14196637909000000000000000000"), cmath.MustParseBig256("27821095556737700000")}, // year 11 {big.NewInt(19582736803651100), cmath.MustParseBig256("14489093265000000000000000000"), cmath.MustParseBig256("26991638121944800000")}, // year 12 {big.NewInt(18800000000000000), cmath.MustParseBig256("14772829365000000000000000000"), cmath.MustParseBig256("26420204724736900000")}, // year 13 } yearInSecs = uint64(31536000) ) var year uint64 // calculate current inflation year from block number year = (currentTime - forkTime) / yearInSecs log.Debug("inflation year", "year", year) if year >= 13 { return inflationData[13][0], inflationData[13][1], inflationData[13][2] } return inflationData[year][0], inflationData[year][1], inflationData[year][2] } func (p *Parlia) getLastSupplyFromTokenomics(header *types.Header) (*big.Int, error) { method := "getTotalSupply" data, err := p.tokenomicsABI.Pack(method) if err != nil { return nil, err } msgData := (hexutil.Bytes)(data) gas := (hexutil.Uint64)(uint64(math.MaxUint64 / 2)) args := ethapi.TransactionArgs{ From: &header.Coinbase, To: &systemcontract.TokenomicsContractAddress, Gas: &gas, Data: &msgData, } // Pin to the parent by hash, never by number (COR-184): by-number resolution goes // through the canonical number->hash index, so a node-local gap in that index // becomes a hard error out of Finalize. RequireCanonical must stay false, so // re-execution off the canonical chain still resolves. blockNr := rpc.BlockNumberOrHashWithHash(header.ParentHash, false) res, err := p.ethAPI.Call(context.Background(), args, &blockNr, nil, nil) if err != nil { return nil, err } initialTotalSupply := big.NewInt(0) if err := p.tokenomicsABI.UnpackIntoInterface(&initialTotalSupply, method, res); err != nil { return nil, err } return initialTotalSupply, nil } func (p *Parlia) distributeToTokenomics(amount *big.Int, inflationPct *big.Int, validator common.Address, newTotalSupply *big.Int, state vm.StateDB, header *types.Header, chain core.ChainContext, txs *[]*types.Transaction, receipts *[]*types.Receipt, receivedTxs *[]*types.Transaction, usedGas *uint64, mode systemTxMode, tracer *tracing.Hooks) error { // method method := "deposit" // get packed data data, err := p.tokenomicsABI.Pack(method, validator, newTotalSupply, inflationPct) if err != nil { log.Error("Unable to pack tx for tokenomics deposit", "error", err) return err } // get system message msg := p.getSystemMessage(header.Coinbase, systemcontract.TokenomicsContractAddress, data, amount) // apply message return p.applyTransaction(msg, state, header, chain, txs, receipts, receivedTxs, usedGas, mode, tracer) } // getValidatorTotalDelegated returns the total delegated amount for a validator at // the epoch containing blockNumber. `state` selects the state the contract call runs // against: nil means latest (pre-Snake8Fix historical behavior — nondeterministic // across nodes); post-Snake8Fix callers pin it to the parent block so producers and // verifiers read identical values (COR-173). func (p *Parlia) getValidatorTotalDelegated(validatorAddress common.Address, blockNumber uint64, state *rpc.BlockNumberOrHash) (*big.Int, error) { if p.ethAPI == nil { // Engines constructed without an RPC backend (tests, tooling) cannot // read stakes; report it as the same soft failure as an RPC error so // frequency selection degrades to round-robin instead of panicking. return nil, errors.New("ethAPI unavailable for stake lookup") } method := "getValidatorStatusAtEpoch" epoch := blockNumber / p.chainConfig.Parlia.Epoch data, err := p.validatorSetABI.Pack(method, validatorAddress, epoch) if err != nil { log.Error("Unable to pack tx for getValidatorStatusAtEpoch", "error", err) return nil, err } ctx, cancel := context.WithCancel(context.Background()) defer cancel() msgData := (hexutil.Bytes)(data) toAddress := common.HexToAddress(systemcontracts.ValidatorContract) gas := (hexutil.Uint64)(uint64(math.MaxUint64 / 2)) result, err := p.ethAPI.Call(ctx, ethapi.TransactionArgs{ Gas: &gas, To: &toAddress, Data: &msgData, }, state, nil, nil) if err != nil { return nil, err } var status struct { OwnerAddress common.Address `json:"ownerAddress"` // Address of the owner Status uint8 `json:"status"` // Status of the validator TotalDelegated *big.Int `json:"totalDelegated"` // Total amount delegated (uint256) SlashesCount uint32 `json:"slashesCount"` // Count of slashes (uint32) ChangedAt uint64 `json:"changedAt"` // Timestamp when status changed (uint64) JailedBefore uint64 `json:"jailedBefore"` // Timestamp when jailed (uint64) ClaimedAt uint64 `json:"claimedAt"` // Timestamp when rewards were claimed (uint64) CommissionRate uint16 `json:"commissionRate"` // Commission rate (uint16) TotalRewards *big.Int `json:"totalRewards"` // Total rewards earned (uint96) } err = p.validatorSetABI.UnpackIntoInterface(&status, method, result) if err != nil { return nil, err } return status.TotalDelegated, nil } func (p *Parlia) distributePRB(state vm.StateDB, header *types.Header, chain core.ChainContext, txs *[]*types.Transaction, receipts *[]*types.Receipt, receivedTxs *[]*types.Transaction, usedGas *uint64, mode systemTxMode, tracer *tracing.Hooks) error { amount := p.GetPepper8MintAmount() recipient := pepper8.Pepper8RecipientAddress log.Info("distributePRB", "amount", amount, "recipient", recipient) // get system message msg := p.getSystemMessage(header.Coinbase, recipient, nil, amount) // apply message return p.applyTransaction(msg, state, header, chain, txs, receipts, receivedTxs, usedGas, mode, tracer) } func (p *Parlia) distributePipe8Mint(state vm.StateDB, header *types.Header, chain core.ChainContext, txs *[]*types.Transaction, receipts *[]*types.Receipt, receivedTxs *[]*types.Transaction, usedGas *uint64, mode systemTxMode, tracer *tracing.Hooks) error { amount := p.GetPipe8MintAmount() recipient := pipe8.Pipe8RecipientAddress log.Info("distributePipe8Mint", "amount", amount, "recipient", recipient) msg := p.getSystemMessage(header.Coinbase, recipient, nil, amount) return p.applyTransaction(msg, state, header, chain, txs, receipts, receivedTxs, usedGas, mode, tracer) } // getCurrentValidators get current validators func (p *Parlia) getCurrentValidators(blockHash common.Hash, blockNum *big.Int) ([]common.Address, map[common.Address]*types.BLSPublicKey, error) { // block blockNr := rpc.BlockNumberOrHashWithHash(blockHash, false) if !p.chainConfig.IsLuban(blockNum) { validators, err := p.getCurrentValidatorsBeforeLuban(blockHash, blockNum) return validators, nil, err } // method method := "getMiningValidators" ctx, cancel := context.WithCancel(context.Background()) defer cancel() // cancel when we are finished consuming integers data, err := p.validatorSetABI.Pack(method) if err != nil { log.Error("Unable to pack tx for getMiningValidators", "error", err) return nil, nil, err } // call msgData := (hexutil.Bytes)(data) toAddress := common.HexToAddress(systemcontract.ValidatorContract) gas := (hexutil.Uint64)(uint64(math.MaxUint64 / 2)) result, err := p.ethAPI.Call(ctx, ethapi.TransactionArgs{ Gas: &gas, To: &toAddress, Data: &msgData, }, &blockNr, nil, nil) if err != nil { return nil, nil, err } var valSet []common.Address var voteAddrSet []types.BLSPublicKey if err := p.validatorSetABI.UnpackIntoInterface(&[]interface{}{&valSet, &voteAddrSet}, method, result); err != nil { return nil, nil, err } voteAddrMap := make(map[common.Address]*types.BLSPublicKey, len(valSet)) for i := 0; i < len(valSet); i++ { voteAddrMap[valSet[i]] = &(voteAddrSet)[i] } return valSet, voteAddrMap, nil } func (p *Parlia) isIntentionalDelayMining(chain consensus.ChainHeaderReader, header *types.Header) (bool, error) { parent := chain.GetHeader(header.ParentHash, header.Number.Uint64()-1) if parent == nil { return false, errors.New("parent not found") } blockInterval, err := p.BlockInterval(chain, header) if err != nil { return false, err } isIntentional := header.Coinbase == parent.Coinbase && header.Difficulty.Cmp(diffInTurn) == 0 && parent.Difficulty.Cmp(diffInTurn) == 0 && parent.MilliTimestamp()+blockInterval < header.MilliTimestamp() return isIntentional, nil } // distributeIncoming distributes system incoming of the block func (p *Parlia) distributeIncoming(val common.Address, state vm.StateDB, header *types.Header, chain core.ChainContext, txs *[]*types.Transaction, receipts *[]*types.Receipt, receivedTxs *[]*types.Transaction, usedGas *uint64, mode systemTxMode, tracer *tracing.Hooks) error { var ( coinbase = header.Coinbase isDragon8 = p.chainConfig.IsDragon8(header.Time) || p.chainConfig.IsDragon8Fix(header.Time) ) parent := chain.GetHeader(header.ParentHash, header.Number.Uint64()-1) if parent == nil { return errors.New("parent not found") } if p.IsPepper8Block(header.Time, parent.Time) { // set the bytecode of the deterministic deployment proxy bytecode, err := hex.DecodeString(p.getPepper8DeterministicDeploymentProxyBytecode()) if err != nil { return err } state.SetCode(p.getPepper8DeterministicDeploymentProxyAddress(), bytecode, tracing.CodeChangeSystemContractUpgrade) // distribute Pepper8 log.Trace("distributePRB", "block hash", header.Number.Uint64()) state.AddBalance(coinbase, uint256.MustFromBig(p.GetPepper8MintAmount()), tracing.BalanceChangeUnspecified) if err := p.distributePRB(state, header, chain, txs, receipts, receivedTxs, usedGas, mode, tracer); err != nil { return err } } if p.IsPipe8Block(header.Time, parent.Time) { log.Trace("distributePipe8Mint", "block hash", header.Number.Uint64()) state.AddBalance(coinbase, uint256.MustFromBig(p.GetPipe8MintAmount()), tracing.BalanceChangeUnspecified) if err := p.distributePipe8Mint(state, header, chain, txs, receipts, receivedTxs, usedGas, mode, tracer); err != nil { return err } } if isDragon8 { var ( blockAmount *big.Int inflationPct *big.Int lastSupply *big.Int newTotalSupply *big.Int ) if p.chainConfig.IsDragon8Fix(header.Time) { inflationPct, newTotalSupply, blockAmount = getNewSupplyForBlockDragon8Fix(*p.chainConfig.Dragon8FixTime, header.Time) } else if p.chainConfig.IsDragon8(header.Time) { // Only this schedule consumes lastSupply, so the Tokenomics read stays // confined to this branch (COR-184) — do not hoist it out. var err error lastSupply, err = p.getLastSupplyFromTokenomics(header) if err != nil { return err } blockAmount, inflationPct = getNewSupplyForBlock(*p.chainConfig.Dragon8Time, header.Time, lastSupply) newTotalSupply = big.NewInt(0).Add(lastSupply, blockAmount) } state.AddBalance(coinbase, uint256.MustFromBig(blockAmount), tracing.BalanceChangeUnspecified) // DEPOSIT to tokenomics // lastSupply is read only on the legacy branch, so it is not logged here. log.Trace("distribute to tokenomics", "block hash", header.Hash(), "amount", blockAmount, "inflation", inflationPct, "newTotalSupply", newTotalSupply) if err := p.distributeToTokenomics(blockAmount, inflationPct, val, newTotalSupply, state, header, chain, txs, receipts, receivedTxs, usedGas, mode, tracer); err != nil { return err } } balance := state.GetBalance(consensus.SystemAddress) if balance.Cmp(common.U2560) <= 0 { return nil } state.SetBalance(consensus.SystemAddress, common.U2560, tracing.BalanceChangeUnspecified) state.AddBalance(coinbase, balance, tracing.BalanceChangeUnspecified) doDistributeSysReward := !isDragon8 && state.GetBalance(common.HexToAddress(systemcontracts.SystemRewardContract)).Cmp(maxSystemBalance) < 0 if doDistributeSysReward { rewards := new(big.Int) rewards = rewards.Div(balance.ToBig(), big.NewInt(systemRewardPercent)) if rewards.Cmp(common.Big0) > 0 { err := p.distributeToSystem(rewards, state, header, chain, txs, receipts, receivedTxs, usedGas, mode, tracer) if err != nil { return err } log.Trace("distribute to system reward pool", "block hash", header.Hash(), "amount", rewards) balance = balance.Sub(balance, uint256.MustFromBig(rewards)) } } log.Trace("distribute to validator contract", "block hash", header.Hash(), "amount", balance) return p.distributeToValidator(balance.ToBig(), val, state, header, chain, txs, receipts, receivedTxs, usedGas, mode, tracer) } // slash spoiled validators func (p *Parlia) slash(spoiledVal common.Address, state vm.StateDB, header *types.Header, chain core.ChainContext, txs *[]*types.Transaction, receipts *[]*types.Receipt, receivedTxs *[]*types.Transaction, usedGas *uint64, mode systemTxMode, tracer *tracing.Hooks) error { // method method := "slash" // get packed data data, err := p.slashABI.Pack(method, spoiledVal, ) if err != nil { log.Error("Unable to pack tx for slash", "error", err) return err } // get system message msg := p.getSystemMessage(header.Coinbase, common.HexToAddress(systemcontract.SlashContract), data, common.Big0) // apply message return p.applyTransaction(msg, state, header, chain, txs, receipts, receivedTxs, usedGas, mode, tracer) } // init contract func (p *Parlia) initContract(state vm.StateDB, header *types.Header, chain core.ChainContext, txs *[]*types.Transaction, receipts *[]*types.Receipt, receivedTxs *[]*types.Transaction, usedGas *uint64, mode systemTxMode, tracer *tracing.Hooks) error { // method method := "init" // get packed data data, err := p.validatorSetABI.Pack(method) if err != nil { log.Error("Unable to pack tx for init validator set", "error", err) return err } contracts := []common.Address{ common.HexToAddress(systemcontract.ValidatorContract), common.HexToAddress(systemcontract.SlashContract), common.HexToAddress(systemcontract.SystemRewardContract), common.HexToAddress(systemcontract.StakingPoolContract), common.HexToAddress(systemcontract.GovernanceContract), common.HexToAddress(systemcontract.ChainConfigContract), common.HexToAddress(systemcontract.RuntimeUpgradeContract), } if !p.chainConfig.IsDeployerProxySunsetTime(header.Time) { contracts = append(contracts, common.HexToAddress(systemcontract.DeployerProxyContract)) } if p.chainConfig.IsDragon8(header.Time) || p.chainConfig.IsDragon8Fix(header.Time) { contracts = append(contracts, common.HexToAddress(systemcontract.TokenomicsContract)) } for _, c := range contracts { msg := p.getSystemMessage(header.Coinbase, c, data, common.Big0) // apply message log.Trace("init contract", "block hash", header.Hash(), "contract", c) err = p.applyTransaction(msg, state, header, chain, txs, receipts, receivedTxs, usedGas, mode, tracer) if err != nil { return err } } return nil } func (p *Parlia) distributeToSystem(amount *big.Int, state vm.StateDB, header *types.Header, chain core.ChainContext, txs *[]*types.Transaction, receipts *[]*types.Receipt, receivedTxs *[]*types.Transaction, usedGas *uint64, mode systemTxMode, tracer *tracing.Hooks) error { // get system message msg := p.getSystemMessage(header.Coinbase, common.HexToAddress(systemcontract.SystemRewardContract), nil, amount) // apply message return p.applyTransaction(msg, state, header, chain, txs, receipts, receivedTxs, usedGas, mode, tracer) } // distributeToValidator deposits validator reward to validator contract func (p *Parlia) distributeToValidator(amount *big.Int, validator common.Address, state vm.StateDB, header *types.Header, chain core.ChainContext, txs *[]*types.Transaction, receipts *[]*types.Receipt, receivedTxs *[]*types.Transaction, usedGas *uint64, mode systemTxMode, tracer *tracing.Hooks) error { // method method := "deposit" // get packed data data, err := p.validatorSetABI.Pack(method, validator, ) if err != nil { log.Error("Unable to pack tx for deposit", "error", err) return err } // get system message msg := p.getSystemMessage(header.Coinbase, common.HexToAddress(systemcontract.ValidatorContract), data, amount) // apply message return p.applyTransaction(msg, state, header, chain, txs, receipts, receivedTxs, usedGas, mode, tracer) } // get system message func (p *Parlia) getSystemMessage(from, toAddress common.Address, data []byte, value *big.Int) *core.Message { return &core.Message{ From: from, GasLimit: math.MaxUint64 / 2, GasPrice: big.NewInt(0), Value: value, To: &toAddress, Data: data, } } func (p *Parlia) applyTransaction( msg *core.Message, state vm.StateDB, header *types.Header, chainContext core.ChainContext, txs *[]*types.Transaction, receipts *[]*types.Receipt, receivedTxs *[]*types.Transaction, usedGas *uint64, mode systemTxMode, tracer *tracing.Hooks, ) (applyErr error) { nonce := state.GetNonce(msg.From) expectedTx := types.NewTransaction(nonce, *msg.To, msg.Value, msg.GasLimit, msg.GasPrice, msg.Data) expectedHash := p.signer.Hash(expectedTx) switch mode { case systemTxMining: var err error if msg.From != p.val { return fmt.Errorf("cannot sign system tx from %s with validator %s", msg.From, p.val) } expectedTx, err = p.signTxFn(accounts.Account{Address: msg.From}, expectedTx, p.chainConfig.ChainID) if err != nil { return err } case systemTxPacking: case systemTxImporting: if receivedTxs == nil || len(*receivedTxs) == 0 || (*receivedTxs)[0] == nil { return errors.New("supposed to get a actual transaction, but get none") } actualTx := (*receivedTxs)[0] if !bytes.Equal(p.signer.Hash(actualTx).Bytes(), expectedHash.Bytes()) { return fmt.Errorf("expected tx hash %v, nonce %d, to %s, value %s, gas %d, gasPrice %s, data %s\ngot tx hash %v, nonce %d, to %s, value %s, gas %d, gasPrice %s, data %s", expectedHash.String(), expectedTx.Nonce(), expectedTx.To().String(), expectedTx.Value().String(), expectedTx.Gas(), expectedTx.GasPrice().String(), hex.EncodeToString(expectedTx.Data()), actualTx.Hash().Hex(), actualTx.Nonce(), actualTx.To().String(), actualTx.Value().String(), actualTx.Gas(), actualTx.GasPrice().String(), hex.EncodeToString(actualTx.Data()), ) } expectedTx = actualTx // move to next *receivedTxs = (*receivedTxs)[1:] default: return fmt.Errorf("unknown system tx mode %d", mode) } state.SetTxContext(expectedTx.Hash(), len(*txs)) // Create a new context to be used in the EVM environment context := core.NewEVMBlockContext(header, chainContext, nil) // Create a new environment which holds all relevant information // about the transaction and calling mechanisms. evm := vm.NewEVM(context, state, p.chainConfig, vm.Config{Tracer: tracer}) evm.SetTxContext(core.NewEVMTxContext(msg)) // Tracing receipt will be set if there is no error and will be used to trace the transaction var tracingReceipt *types.Receipt if tracer != nil { if tracer.OnSystemTxStart != nil { tracer.OnSystemTxStart() } if tracer.OnTxStart != nil { tracer.OnTxStart(evm.GetVMContext(), expectedTx, msg.From) } // Defers are last in first out, so OnTxEnd will run before OnSystemTxEnd in this transaction, // which is what we want. if tracer.OnSystemTxEnd != nil { defer func() { tracer.OnSystemTxEnd() }() } if tracer.OnTxEnd != nil { defer func() { tracer.OnTxEnd(tracingReceipt, applyErr) }() } } gasUsed, err := applyMessage(msg, evm, state, header, p.chainConfig, chainContext) if err != nil { return err } *txs = append(*txs, expectedTx) // increment nonce only when tx is included state.SetNonce(msg.From, nonce+1, tracing.NonceChangeEoACall) var root []byte if p.chainConfig.IsByzantium(header.Number) { state.Finalise(true) } else { root = state.IntermediateRoot(p.chainConfig.IsEIP158(header.Number)).Bytes() } *usedGas += gasUsed tracingReceipt = types.NewReceipt(root, false, *usedGas) tracingReceipt.TxHash = expectedTx.Hash() tracingReceipt.GasUsed = gasUsed // Set the receipt logs and create a bloom for filtering tracingReceipt.Logs = state.GetLogs(expectedTx.Hash(), header.Number.Uint64(), header.Hash(), header.Time) tracingReceipt.Bloom = types.CreateBloom(tracingReceipt) tracingReceipt.BlockHash = header.Hash() tracingReceipt.BlockNumber = header.Number tracingReceipt.TransactionIndex = uint(state.TxIndex()) *receipts = append(*receipts, tracingReceipt) return nil } // GetJustifiedNumberAndHash retrieves the number and hash of the highest justified block // within the branch including `headers` and utilizing the latest element as the head. func (p *Parlia) GetJustifiedNumberAndHash(chain consensus.ChainHeaderReader, headers []*types.Header) (uint64, common.Hash, error) { if chain == nil || len(headers) == 0 || headers[len(headers)-1] == nil { return 0, common.Hash{}, errors.New("illegal chain or header") } head := headers[len(headers)-1] snap, err := p.snapshot(chain, head.Number.Uint64(), head.Hash(), headers, p.isSnake8Enabled(chain, head), head) if err != nil { log.Error("GJ Unexpected error when getting snapshot", "error", err, "blockNumber", head.Number.Uint64(), "blockHash", head.Hash()) return 0, common.Hash{}, err } if snap.Attestation == nil { if p.chainConfig.IsLuban(head.Number) { log.Debug("once one attestation generated, attestation of snap would not be nil forever basically") } return 0, chain.GetHeaderByNumber(0).Hash(), nil } return snap.Attestation.TargetNumber, snap.Attestation.TargetHash, nil } // GetFinalizedHeader returns highest finalized block header. // It first checks VotePool for votes that may have reached quorum but not yet included in block headers, // then falls back to the attestation in the snapshot. func (p *Parlia) GetFinalizedHeader(chain consensus.ChainHeaderReader, header *types.Header) *types.Header { if chain == nil || header == nil { return nil } if !chain.Config().IsPlato(header.Number) { return chain.GetHeaderByNumber(0) } parent := chain.GetHeaderByHash(header.ParentHash) if parent == nil { log.Error("parent not found") return nil } snap, err := p.snapshot(chain, header.Number.Uint64(), header.Hash(), nil, p.isSnake8Enabled(chain, header), header) if err != nil { log.Error("GF Unexpected error when getting snapshot", "error", err, "blockNumber", header.Number.Uint64(), "blockHash", header.Hash()) return nil } if snap.Attestation == nil { return chain.GetHeaderByNumber(0) // keep consistent with GetJustifiedNumberAndHash } finalizedHash := snap.Attestation.SourceHash finalizedNumber := snap.Attestation.SourceNumber currentJustifiedHash := snap.Attestation.TargetHash currentJustifiedNumber := snap.Attestation.TargetNumber // Try to check if currentJustifiedNumber can become finalized by checking VotePool. // We only need to check currentJustifiedNumber + 1, since currentJustifiedNumber is already the latest justified. if p.VotePool != nil && currentJustifiedNumber == header.Number.Uint64()-1 { parentSnap, err := p.snapshot(chain, header.Number.Uint64()-1, header.ParentHash, nil, p.isSnake8Enabled(chain, header), header) if err == nil { // Check if the next block (direct child) has reached quorum in VotePool votes := p.VotePool.FetchVotesByBlockHash(header.Hash(), currentJustifiedNumber) quorum := cmath.CeilDiv(len(parentSnap.Validators)*2, 3) if len(votes) >= quorum { finalizedHash = currentJustifiedHash finalizedNumber = currentJustifiedNumber } } else { log.Error("Failed to get parent snapshot for finality check", "error", err, "blockNumber", header.Number.Uint64()-1, "blockHash", header.ParentHash) } } return chain.GetHeader(finalizedHash, finalizedNumber) } // CheckFinalityAndNotify checks if votes for the target block have reached quorum, // and if so, notifies the blockchain of early finalization via the notifyFn callback. func (p *Parlia) CheckFinalityAndNotify(chain consensus.ChainHeaderReader, targetBlockHash common.Hash, notifyFn func(finalizedHeader *types.Header)) { // Get target block header directly by hash (don't rely on currentHeader which may have moved forward) targetHeader := chain.GetHeaderByHash(targetBlockHash) if targetHeader == nil { return } finalizedHeader := p.GetFinalizedHeader(chain, targetHeader) if finalizedHeader == nil || finalizedHeader.Number.Uint64() == 0 { return } // Notify via callback (NotifyFinalized has its own deduplication logic) notifyFn(finalizedHeader) } // =========================== utility function ========================== func (p *Parlia) backOffTime(snap *Snapshot, parent, header *types.Header, val common.Address) uint64 { if snap.inturn(val) { log.Debug("backOffTime", "blockNumber", header.Number, "in turn validator", val) return 0 } else { delay := defaultInitialBackOffTime // When mining blocks, `header.Time` is temporarily set to time.Now() + 1. // Therefore, using `header.Time` to determine whether a hard fork has occurred is incorrect. // As a result, during the Bohr and Lorentz hard forks, the network may experience some instability, // So use `parent.Time` instead. isParerntLorentz := p.chainConfig.IsLorentz(parent.Number, parent.Time) if isParerntLorentz { // If the in-turn validator has not signed recently, the expected backoff times are [2, 3, 4, ...]. delay = lorentzInitialBackOffTime } validators := snap.validators() if p.chainConfig.IsPlanck(header.Number) { counts := snap.countRecents() for addr, seenTimes := range counts { log.Trace("backOffTime", "blockNumber", header.Number, "validator", addr, "seenTimes", seenTimes) } // The backOffTime does not matter when a validator has signed recently. if snap.signRecentlyByCounts(val, counts) { return 0 } inTurnAddr := snap.inturnValidator() if snap.signRecentlyByCounts(inTurnAddr, counts) { log.Debug("in turn validator has recently signed, skip initialBackOffTime", "inTurnAddr", inTurnAddr) delay = 0 } // Exclude the recently signed validators and the in turn validator temp := make([]common.Address, 0, len(validators)) for _, addr := range validators { if snap.signRecentlyByCounts(addr, counts) { continue } if p.chainConfig.IsBohr(header.Number, header.Time) { if addr == inTurnAddr { continue } } temp = append(temp, addr) } validators = temp } // get the index of current validator and its shuffled backoff time. idx := -1 for index, itemAddr := range validators { if val == itemAddr { idx = index } } if idx < 0 { log.Debug("The validator is not authorized", "addr", val) return 0 } randSeed := snap.Number if p.chainConfig.IsBohr(header.Number, header.Time) { randSeed = header.Number.Uint64() / uint64(snap.TurnLength) } s := rand.NewSource(int64(randSeed)) r := rand.New(s) n := len(validators) backOffSteps := make([]uint64, 0, n) for i := uint64(0); i < uint64(n); i++ { backOffSteps = append(backOffSteps, i) } r.Shuffle(n, func(i, j int) { backOffSteps[i], backOffSteps[j] = backOffSteps[j], backOffSteps[i] }) if delay == 0 && isParerntLorentz { // If the in-turn validator has signed recently, the expected backoff times are [0, 2, 3, ...]. if backOffSteps[idx] == 0 { return 0 } return lorentzInitialBackOffTime + (backOffSteps[idx]-1)*wiggleTime } delay += backOffSteps[idx] * wiggleTime return delay } } // BlockInterval returns number of blocks in one epoch for the given header func (p *Parlia) epochLength(chain consensus.ChainHeaderReader, header *types.Header, parents []*types.Header) (uint64, error) { if header == nil { return defaultEpochLength, errUnknownBlock } if header.Number.Uint64() == 0 { return defaultEpochLength, nil } // extraHeader is nil: this caller only reads EpochLength and has no candidate // header whose frequency data it would want stamped onto its snapshot view. snap, err := p.snapshot(chain, header.Number.Uint64()-1, header.ParentHash, parents, p.isSnake8Enabled(chain, header), nil) if err != nil { return defaultEpochLength, err } return snap.EpochLength, nil } // BlockInterval returns the block interval in milliseconds for the given header func (p *Parlia) BlockInterval(chain consensus.ChainHeaderReader, header *types.Header) (uint64, error) { if header == nil { return defaultBlockInterval, errUnknownBlock } if header.Number.Uint64() == 0 { return defaultBlockInterval, nil } snap, err := p.snapshot(chain, header.Number.Uint64()-1, header.ParentHash, nil, p.isSnake8Enabled(chain, header), nil) if err != nil { return defaultBlockInterval, err } return snap.BlockInterval, nil } func (p *Parlia) NextProposalBlock(chain consensus.ChainHeaderReader, header *types.Header, proposer common.Address) (uint64, uint64, error) { snap, err := p.snapshot(chain, header.Number.Uint64(), header.Hash(), nil, p.isSnake8Enabled(chain, header), header) if err != nil { return 0, 0, err } return snap.nextProposalBlock(proposer) } func (p *Parlia) detectNewVersionWithFork(chain consensus.ChainHeaderReader, header *types.Header, state vm.StateDB) { // Ignore blocks that are considered too old const maxBlockReceiveDelay = 10 * time.Second blockTime := time.UnixMilli(int64(header.MilliTimestamp())) if time.Since(blockTime) > maxBlockReceiveDelay { return } // If the fork is not a majority, log a warning or debug message number := header.Number.Uint64() snap, err := p.snapshot(chain, number-1, header.ParentHash, nil, p.isSnake8Enabled(chain, header), header) if err != nil { return } nextForkHash := forkid.NextForkHash(p.chainConfig, p.genesisHash, chain.GenesisHeader().Time, number, header.Time) forkHashHex := hex.EncodeToString(nextForkHash[:]) if !snap.isMajorityFork(forkHashHex) { logFn := log.Debug if state.NoTries() { logFn = log.Warn } logFn("possible fork detected: client is not in majority", "nextForkHash", forkHashHex) } } // TODO(Nathan): use kAncestorGenerationDepth directly instead of this func once Fermi hardfork passed func (p *Parlia) GetAncestorGenerationDepth(header *types.Header) uint64 { if p.chainConfig.IsFermi(header.Number, header.Time) { return kAncestorGenerationDepth } return 1 } // chain context type chainContext struct { consensus.ChainHeaderReader parlia consensus.Engine } func (c chainContext) Engine() consensus.Engine { return c.parlia } // apply message func applyMessage( msg *core.Message, evm *vm.EVM, state vm.StateDB, header *types.Header, chainConfig *params.ChainConfig, chainContext core.ChainContext, ) (uint64, error) { // Apply the transaction to the current state (included in the env) if chainConfig.IsCancun(header.Number, header.Time) { rules := evm.ChainConfig().Rules(evm.Context.BlockNumber, evm.Context.Random != nil, evm.Context.Time) state.Prepare(rules, msg.From, evm.Context.Coinbase, msg.To, vm.ActivePrecompiles(rules), msg.AccessList) } else { state.ClearAccessList() } ret, returnGas, err := evm.Call( msg.From, *msg.To, msg.Data, msg.GasLimit, uint256.MustFromBig(msg.Value), ) if err != nil && len(ret) > 64+4 { log.Error("apply message failed", "msg", string(ret[64+4:]), "err", err) } return msg.GasLimit - returnGas, err } // proposalKey build a key which is a combination of the block number and the proposer address. func proposalKey(header types.Header) string { return header.ParentHash.String() + header.Coinbase.String() }