From 4af39c8a5c6f23ef0619743f290c1c541c8c7acb Mon Sep 17 00:00:00 2001
From: Daniel Liu <139250065@qq.com>
Date: Tue, 8 Sep 2026 22:31:08 +0800
Subject: [PATCH] fix(eth,consensus): gate the proposed-block handler on
canonicality and storage
MIME-Version: 1.0
Content-Type: text/plain; charset=UTF-8
Content-Transfer-Encoding: 8bit
importBlockResults handed the tail of every batch to the proposed-block
handler as long as InsertChain returned nil. A nil error does not mean
the tail is canonical: a fork batch is written as side-chain entries, and
a parked tail is not written at all. The engine did not make up for it:
processQC updates highestQuorumCert, lockQuorumCert and the commit block
before its own existence check, and an existence check cannot tell a
reorged-away block from a canonical one, because the fork stays in the
database. The fetcher and the miner call the same handler and had no gate
at all. Canonicality is not enough on its own either: it is a property of
the header, and the fast sync header phase marks a height canonical
before its body lands, so a node could processQC and vote for a block
that only ever existed as a header. A master node could therefore end up
voting for, and committing state against, a block it had just reorged
away.
Judge the block once, in one place. consensus.ShouldHandleProposedBlock
reports whether a header is the canonical block at its height and whether
its body is stored, together with the reason it must be skipped and the
canonical hash at that height, both for the skip log. It takes a minimal
CanonicalChain (GetHeaderByNumber alone) that both consensus.ChainReader
and the downloader's BlockChain satisfy, so the callers cannot drift into
two diverging judgments; on the downloader side that costs one addition of
GetHeaderByNumber to its BlockChain interface. The storage half goes
through the optional capability interface consensus.BlockStorer, whose
single method is HasBlock rather than GetBlock: only existence matters,
and GetBlock would read and RLP-decode the whole body of every imported
block. core.BlockChain implements BlockStorer — a compile-time assertion
keeps it from dropping off the type — while a HeaderChain deliberately
does not, so handing a header-only chain to the judgment fails it loudly
as the SkipUnjudgeable skip reason instead of every block silently
failing as not stored. The judgment has no error channel: every outcome
is a skip reason, so each caller collapses to a single !ok branch. The
BlockStorer assertion runs before any chain read, so the unjudgeable
skip is reported at every height: a header-only chain must not get a
silent "not stored" from the canonicality half for the heights it
happens to answer.
The downloader gates importBlockResults on that judgment instead of
open-coding the same rule out of HasBlock and GetCanonicalHash, and logs
the skip with its reason and the canonical hash it observed. The v2 engine
re-checks the same judgment twice inside ProposedBlockHandler: in front of
processQC, and again right before sendVote, because x.lock serialises the
handler but not InsertChain, so a reorg can still land between the two
checks. A judgment that cannot run — a chain type that lost the
BlockStorer half to wrapping or replacement — surfaces as the
SkipUnjudgeable skip, which the handler logs and returns as nil: the
fetcher's import loop treats any handler error as an import failure and
would suppress the broadcast of an already imported block, so a skip
must never surface as an error. That skip means the node stops
processing QCs and voting outright — in production the chain always
implements BlockStorer, so it is a wiring bug — so the handler also
increments a metrics counter (consensus/enginev2/unjudgeable) to make
the liveness halt observable beyond the logs.
consensus.SkipLogLevel grades the skips by reason:
SkipUnjudgeable is the one Error — not an observation about the block but
a wiring bug that will skip every block — SkipNonCanonical is a genuine
reorg race and stays at Warn, while the sync-phase skips are routine and
stay at Info so they do not drown the level in noise. The downloader's
own skip log stays at Info on purpose: it is the pre-filter for the
routine cases and correctness rests on the handler's re-checks, so
grading it too would fire a Warn for every fork tail of a sync; its
unjudgeable branch is an explicit Error instead of SkipLogLevel.
The fetcher is the one caller whose fast sync calls must not reach the
handler at all. While snapSync runs, it discards propagated blocks before
executing them, so a body already written by the fast sync receipt phase
would pass both halves of the judgment — canonical and stored — and drive
processQC and the vote path on a block whose state transition was never
validated. Its callback therefore carries the same snapSync guard the
inserter and prepare closures already have, instead of asking the
judgment a question the judgment cannot answer. The guard sits in a
fetcher-specific wrapper and not in the shared closure: the same closure
also feeds the downloader, whose fast sync handler calls run after the
pivot commit, on blocks InsertChain has fully executed, and snapSync
stays set until the whole Synchronise returns. Behind that flag the
callback also requires the block's full state, keyed by the block hash
(HasBlockAndFullState) rather than the state root: an empty block's root
can repeat its parent's, so a root-keyed check would leak the parent's
executed state to a block whose import was discarded. The fetcher's skip
is logged at Debug, so a syncing node can still see why no QC processing
happens without drowning the level.
Tests:
- TestShouldHandleProposedBlock covers the five outcomes of the shared
judgment plus the header-only chain at an absent height, whose
unjudgeable skip is the point of the entry-side assertion, and
TestHeaderChainShouldHandleProposedBlock pins that the header chain's
interface stubs never pass it at any height.
- TestProposedBlockHandlerSkipsNonCanonicalBlock,
TestProposedBlockHandlerSkipsReorgedBlockBeforeProcessQC,
TestProposedBlockHandlerDropsVoteForReorgedBlock and
TestProposedBlockHandlerSkipsBlockWithoutBody cover the two engine
re-checks: highestQuorumCert, lockQC, the timeout certificate, the
voted round and the commit block must stay untouched and no vote may be
broadcast. The two reorg tests inject the reorg through a ChainReader
wrapper that serves a fork header after N truthful reads, and assert
the read count, so a chain read added to the handler cannot silently
move where the injection lands.
- TestProposedBlockHandlerGradesSkipLogLevelByReason pins the log level
of each skip reason.
- TestProposedBlockHandlerSkipsUnjudgeableChain feeds the handler a
ChainReader wrapper without the BlockStorer half — the shape a wrapped
or replaced chain type presents — and pins that the chain skips as
SkipUnjudgeable, logged at Error and returned as nil, not surfaced to a
caller that would treat it as an import failure.
- TestFetcherSkipsProposedBlockHandlerDuringFastSync pins the fetcher
gate: with snapSync set the fetcher's callback is inert and logs the
skip, and the same call in full sync reaches the closure without the
skip log. It also reads the handler back off the fetcher through the
test-only HandleProposedBlock accessor and requires it to be the gated
closure itself, so swapping the wiring back to the bare
handleProposedBlock fails loudly instead of staying green.
- TestImportBlockResultsProposedBlockHandler covers six downloader
shapes: a parked tail, a fully imported batch, a stored fork batch
re-delivered after the local chain grew past it, a head advanced past
the canonical tail by a concurrent import, a heavier fork that stays
canonical, and a fast sync height that is canonical without a body.
- downloadTester gains a canonical number-to-hash table, picked by total
difficulty the way the real chain resolves a reorg, plus hooks to park
a batch tail and to extend it after the insert.
TestStoreBlockCleansStaleCanonicalMarkers,
TestRollbackClearsCanonicalMarkers and
TestInsertChainErrorReportsPosition pin the parts of that table the
downloader tests rely on.
- TestShouldNotSendVoteMsgIfBlockNotExtendedFromAncestor no longer
proposes a forked block, which the new entry re-check now
short-circuits; it proposes a canonical block below the locked ancestor
instead. The fork case moved to
TestShouldNotSendVoteMsgIfCanonicalBlockNotExtendedFromForkedAncestor,
where the parent walk of isExtendingFromAncestor actually runs.
---
consensus/XDPoS/engines/engine_v2/engine.go | 46 ++
.../XDPoS/engines/engine_v2/vote_test.go | 60 ++
consensus/proposed_block.go | 151 ++++
consensus/proposed_block_test.go | 212 +++++
.../engine_v2_tests/proposed_block_test.go | 498 +++++++++++-
core/blockchain.go | 9 +
core/headerchain.go | 11 +-
core/headerchain_test.go | 50 ++
eth/downloader/downloader.go | 35 +-
eth/downloader/downloader_test.go | 751 +++++++++++++++++-
eth/fetcher/block_fetcher.go | 9 +
eth/handler.go | 38 +-
eth/handler_test.go | 84 ++
13 files changed, 1911 insertions(+), 43 deletions(-)
create mode 100644 consensus/proposed_block.go
create mode 100644 consensus/proposed_block_test.go
create mode 100644 core/headerchain_test.go
diff --git a/consensus/XDPoS/engines/engine_v2/engine.go b/consensus/XDPoS/engines/engine_v2/engine.go
index c4c86684b989..8512817f6675 100644
--- a/consensus/XDPoS/engines/engine_v2/engine.go
+++ b/consensus/XDPoS/engines/engine_v2/engine.go
@@ -27,11 +27,19 @@ import (
"github.com/XinFinOrg/XDPoSChain/core/vm"
"github.com/XinFinOrg/XDPoSChain/ethdb"
"github.com/XinFinOrg/XDPoSChain/log"
+ "github.com/XinFinOrg/XDPoSChain/metrics"
"github.com/XinFinOrg/XDPoSChain/params"
"github.com/XinFinOrg/XDPoSChain/trie"
"golang.org/x/sync/errgroup"
)
+// unjudgeableProposedBlock counts the proposed-block handler skips caused by
+// a chain whose type hides the BlockStorer half of the proposed-block
+// judgment (consensus.SkipUnjudgeable). In production the chain always
+// implements the interface, so a non-zero value is a wiring bug that halts QC
+// processing and voting — a liveness incident that logs alone could bury.
+var unjudgeableProposedBlock = metrics.NewRegisteredCounter("consensus/enginev2/unjudgeable", nil)
+
type XDPoS_v2 struct {
chainConfig *params.ChainConfig // Chain & network configuration
@@ -834,6 +842,36 @@ func (x *XDPoS_v2) ProposedBlockHandler(chain consensus.ChainReader, blockHeader
return err
}
+ // Re-check canonicality and storage at two points. x.lock only serializes
+ // this handler, while the chain is written under the import lock of
+ // BlockChain.InsertChain, so the callers' own gates (downloader,
+ // fetcher, miner) cannot make these checks atomic with a write: a
+ // concurrent reorg can land anywhere in between. The first check
+ // directly protects processQC, which updates highestQuorumCert,
+ // lockQuorumCert and the commit block before its own existence check;
+ // the second check right before sendVote keeps the unguarded window of
+ // the vote down to the broadcast itself. An existence check cannot tell
+ // a reorged-away block from a canonical one, since a fork stays in the
+ // database as side entries.
+ ok, reason, canonicalHash := consensus.ShouldHandleProposedBlock(chain, blockHeader)
+ if !ok {
+ consensus.SkipLogLevel(reason)("[ProposedBlockHandler] skip block before processQC", "reason", reason, "hash", blockHeader.Hash(), "number", blockHeader.Number, "canonicalHash", canonicalHash)
+ if reason == consensus.SkipUnjudgeable {
+ // This skip is the node stopping voting: with the BlockStorer half
+ // off the chain type, every proposed block lands here and
+ // processQC and sendVote never run again. In production the chain
+ // always implements BlockStorer, so this is a wiring bug (a
+ // wrapped or replaced chain type) most likely introduced by a
+ // refactor — the counter makes the liveness halt observable
+ // beyond the logs.
+ unjudgeableProposedBlock.Inc(1)
+ }
+ // Return nil, never an error: the fetcher's import loop treats any
+ // handler error as an import failure and suppresses the broadcast
+ // of an already imported block.
+ return nil
+ }
+
// Generate blockInfo
blockInfo := &types.BlockInfo{
Hash: blockHeader.Hash(),
@@ -856,6 +894,14 @@ func (x *XDPoS_v2) ProposedBlockHandler(chain consensus.ChainReader, blockHeader
return err
}
if verified {
+ // x.lock does not block InsertChain, so a reorg can still land
+ // between the processQC re-check and the broadcast; drop the vote if
+ // the block has been reorged away meanwhile.
+ ok, reason, canonicalHash = consensus.ShouldHandleProposedBlock(chain, blockHeader)
+ if !ok {
+ consensus.SkipLogLevel(reason)("[ProposedBlockHandler] skip vote for reorged block", "reason", reason, "hash", blockHeader.Hash(), "number", blockHeader.Number, "canonicalHash", canonicalHash)
+ return nil
+ }
return x.sendVote(chain, blockInfo)
}
diff --git a/consensus/XDPoS/engines/engine_v2/vote_test.go b/consensus/XDPoS/engines/engine_v2/vote_test.go
index d8c93db79205..ffd12a195761 100644
--- a/consensus/XDPoS/engines/engine_v2/vote_test.go
+++ b/consensus/XDPoS/engines/engine_v2/vote_test.go
@@ -122,3 +122,63 @@ func TestVerifyVoteMessage_VoteRoundTooOld(t *testing.T) {
assert.False(t, verified, "Should return false for vote with round < currentRound")
assert.NoError(t, err, "Should not return an error for old round votes")
}
+
+// blockInfoOf turns a header into the BlockInfo shape the voting rule and
+// forensics pass around. The round is irrelevant to isExtendingFromAncestor,
+// which only walks hashes and numbers.
+func blockInfoOf(h *types.Header) *types.BlockInfo {
+ return &types.BlockInfo{Hash: h.Hash(), Number: h.Number}
+}
+
+// TestIsExtendingFromAncestor covers the parent walk of the HotStuff voting
+// rule at the rule layer. The handler-level tests cannot reach the positive
+// branch anymore: since ProposedBlockHandler gates on canonicality, a
+// proposed block on the locked ancestor's own chain that passes the gate
+// always outranks the lockQC round and returns before the walk (see
+// TestShouldNotSendVoteMsgIfCanonicalBlockNotExtendedFromForkedAncestor),
+// and the forensics caller's positive path is only exercised by a skipped
+// test. Both branches of the walk are safety-critical — a false positive
+// lets a node vote off the locked chain, a false negative stalls it — so
+// the walk itself gets direct coverage here.
+func TestIsExtendingFromAncestor(t *testing.T) {
+ // 1 <- 2 <- 3, with 2' a same-height fork of 2.
+ mockChain := NewMockChainReader()
+ h1 := &types.Header{Number: big.NewInt(1)}
+ h2 := &types.Header{Number: big.NewInt(2), ParentHash: h1.Hash()}
+ h3 := &types.Header{Number: big.NewInt(3), ParentHash: h2.Hash()}
+ forkH2 := &types.Header{Number: big.NewInt(2), ParentHash: h1.Hash(), Coinbase: common.BytesToAddress([]byte{0x02})}
+ for _, h := range []*types.Header{h1, h2, h3} {
+ mockChain.AddHeader(h)
+ }
+ engine := &XDPoS_v2{}
+
+ // Positive branch: the walk runs two hops down the parent chain and
+ // lands exactly on the locked ancestor.
+ extended, err := engine.isExtendingFromAncestor(mockChain, blockInfoOf(h3), blockInfoOf(h1))
+ assert.NoError(t, err)
+ assert.True(t, extended, "h3 extends the locked ancestor h1")
+
+ // Negative branch with the walk executed: h3's parent chain bottoms out
+ // at h1, not at the forked ancestor 2', so the final hash comparison
+ // rejects the block. This is the geometry
+ // TestShouldNotSendVoteMsgIfCanonicalBlockNotExtendedFromForkedAncestor
+ // exercises through verifyVotingRule.
+ extended, err = engine.isExtendingFromAncestor(mockChain, blockInfoOf(h3), blockInfoOf(forkH2))
+ assert.NoError(t, err)
+ assert.False(t, extended, "h3 does not extend the forked ancestor 2'")
+
+ // Zero-iteration mismatch: the proposed block sits below the locked
+ // ancestor, so the walk never runs and only the direct hash comparison
+ // can reject it. This is the geometry of
+ // TestShouldNotSendVoteMsgIfBlockNotExtendedFromAncestor.
+ extended, err = engine.isExtendingFromAncestor(mockChain, blockInfoOf(h1), blockInfoOf(h2))
+ assert.NoError(t, err)
+ assert.False(t, extended, "h1 is below the locked ancestor h2")
+
+ // A missing parent aborts the walk with an error instead of silently
+ // reporting false: the proposed block's own header is not in the chain.
+ missing := &types.BlockInfo{Hash: common.StringToHash("missing"), Number: big.NewInt(3)}
+ extended, err = engine.isExtendingFromAncestor(mockChain, missing, blockInfoOf(h1))
+ assert.Error(t, err)
+ assert.False(t, extended)
+}
diff --git a/consensus/proposed_block.go b/consensus/proposed_block.go
new file mode 100644
index 000000000000..b0006fc0cdfa
--- /dev/null
+++ b/consensus/proposed_block.go
@@ -0,0 +1,151 @@
+// Copyright 2026 The go-ethereum Authors
+// This file is part of the go-ethereum library.
+//
+// The go-ethereum library is free software: you can redistribute it and/or modify
+// it under the terms of the GNU Lesser General Public License as published by
+// the Free Software Foundation, either version 3 of the License, or
+// (at your option) any later version.
+//
+// The go-ethereum library is distributed in the hope that it will be useful,
+// but WITHOUT ANY WARRANTY; without even the implied warranty of
+// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
+// GNU Lesser General Public License for more details.
+//
+// You should have received a copy of the GNU Lesser General Public License
+// along with the go-ethereum library. If not, see .
+
+package consensus
+
+import (
+ "github.com/XinFinOrg/XDPoSChain/common"
+ "github.com/XinFinOrg/XDPoSChain/core/types"
+ "github.com/XinFinOrg/XDPoSChain/log"
+)
+
+// CanonicalChain is the subset of ChainReader needed for the canonicality half
+// of the judgment whether a proposed block should be handled by the consensus
+// engine.
+type CanonicalChain interface {
+ // GetHeaderByNumber retrieves a block header from the database by number.
+ GetHeaderByNumber(number uint64) *types.Header
+}
+
+// BlockStorer is the optional capability the judgment needs for its storage
+// half: whether a block is present in the database. Chains that only carry
+// headers (e.g. *core.HeaderChain) deliberately do not implement it — handing
+// one to ShouldHandleProposedBlock surfaces as a SkipUnjudgeable skip, graded
+// Error, instead of every block silently failing as "not stored".
+type BlockStorer interface {
+ // HasBlock reports whether a block with the given hash and number is
+ // stored in the database.
+ HasBlock(hash common.Hash, number uint64) bool
+}
+
+// SkipReason describes why a proposed block header was rejected by
+// ShouldHandleProposedBlock. The empty value means the header was accepted.
+type SkipReason string
+
+// The reasons ShouldHandleProposedBlock can report for rejecting a header.
+const (
+ // SkipNoCanonicalHeader: no canonical header exists at the height.
+ SkipNoCanonicalHeader SkipReason = "no canonical header at height"
+
+ // SkipNonCanonical: the header is not the canonical block at its height.
+ SkipNonCanonical SkipReason = "non-canonical"
+
+ // SkipBodyNotStored: the canonical block's body has not landed in the
+ // database yet, e.g. during the fast sync header phase.
+ SkipBodyNotStored SkipReason = "block body not stored"
+
+ // SkipUnjudgeable: the chain does not implement BlockStorer, so the
+ // judgment cannot run at all. It is not a judgment about the block but
+ // about the chain: in production the node's chain always implements the
+ // interface, so this reason means a wiring bug — a wrapped or replaced
+ // chain type — and the node will skip every block, stopping processQC
+ // and voting outright. That liveness halt is why it is graded Error.
+ SkipUnjudgeable SkipReason = "unjudgeable"
+)
+
+// ShouldHandleProposedBlock reports whether a proposed block header should
+// reach the consensus handler, together with the reason it should not and the
+// canonical hash at its height, both meant for skip logging. The header must
+// be the canonical block at its height — a fork stays in the database as a
+// side entry, so a mere existence check cannot tell a reorged-away block from
+// a canonical one — and its body must be stored, since the fast sync header
+// phase marks a height canonical before its body lands. HasBlock keeps the
+// storage check off the full-block read path: GetBlock would pull and decode
+// the whole body on every judge call, while only its existence matters here.
+// The skip reason is one of the exported SkipReason constants; the empty
+// reason means the header was accepted.
+// The two chain reads are not atomic: a concurrent reorg can land between
+// GetHeaderByNumber and the storage check, and again between a caller's
+// checks and its effects on chain state. The v2 engine re-checks before
+// processQC and before sendVote; what window remains is recorded there.
+// A chain that does not implement BlockStorer cannot be judged at all:
+// rather than guessing and silently dropping QC processing and voting for
+// every block, the judgment fails it loudly as SkipUnjudgeable. The
+// BlockStorer assertion runs before any chain read, so the skip is
+// reported at every height — a header-only chain must not get a silent
+// "not stored" from the canonicality half for heights it happens to answer.
+func ShouldHandleProposedBlock(chain CanonicalChain, header *types.Header) (bool, SkipReason, common.Hash) {
+ storer, ok := chain.(BlockStorer)
+ if !ok {
+ return false, SkipUnjudgeable, common.Hash{}
+ }
+ canonical := chain.GetHeaderByNumber(header.Number.Uint64())
+ if canonical == nil {
+ return false, SkipNoCanonicalHeader, common.Hash{}
+ }
+ if canonical.Hash() != header.Hash() {
+ return false, SkipNonCanonical, canonical.Hash()
+ }
+ if !storer.HasBlock(header.Hash(), header.Number.Uint64()) {
+ return false, SkipBodyNotStored, canonical.Hash()
+ }
+ return true, "", canonical.Hash()
+}
+
+// SkipLogLevel maps a skip reason to the level its skip should be logged at.
+// Grade the skip by what the reason means once the handler is reached.
+// SkipUnjudgeable is the one Error: it is not a judgment about the block
+// but a wiring bug that will skip every block and halt QC processing and
+// voting, so it must surface beyond the Info/Warn discipline below.
+// SkipNonCanonical means another block already claims this height, and
+// SkipNoCanonicalHeader means no canonical marker exists at all: the fast
+// sync header phase marks heights canonical, so a marker can only be
+// missing post-sync — a fork growing above the local head, a reorged-away
+// tip being re-delivered, or a reorg racing the handler between its
+// checks. Both are reorg-race observations a Warn is reserved for.
+// SkipBodyNotStored is the one routine skip instead: the fast sync header
+// phase marks a height canonical before its body lands, and a Warn per
+// header-only height would drown the level in noise. Info for the
+// expected, Warn for the anomalous — the same discipline the downloader's
+// expected fork-tail skip already follows. The empty reason means the
+// header was accepted and should never reach a skip log; a caller passing
+// it has misused this helper, so it grades like the routine Info skips
+// rather than inflating the Warn level reserved for anomalies. Any other
+// unregistered reason is graded Warn on purpose: a skip reason that was
+// never classified is most likely a future anomalous one whose registration
+// was forgotten, and an Info would bury it exactly where it hurts most.
+// SkipUnjudgeable, finally, is graded Error and sits outside the Info/Warn
+// discipline: it is not an observation about the block at all.
+func SkipLogLevel(reason SkipReason) func(msg string, ctx ...interface{}) {
+ switch reason {
+ case SkipUnjudgeable:
+ // Not an observation about the block but a wiring bug that skips
+ // every block and halts QC processing and voting — a liveness
+ // incident that must be surfaced at Error.
+ return log.Error
+ case SkipNonCanonical, SkipNoCanonicalHeader:
+ return log.Warn
+ case "":
+ // The empty reason means "accepted" (see SkipReason); grade a
+ // misused accept like the routine skips, not like the anomalies.
+ return log.Info
+ case SkipBodyNotStored:
+ return log.Info
+ default:
+ // Unregistered reason: grade it as the anomaly it most likely is.
+ return log.Warn
+ }
+}
diff --git a/consensus/proposed_block_test.go b/consensus/proposed_block_test.go
new file mode 100644
index 000000000000..7d606e5fea9a
--- /dev/null
+++ b/consensus/proposed_block_test.go
@@ -0,0 +1,212 @@
+// Copyright 2026 The go-ethereum Authors
+// This file is part of the go-ethereum library.
+//
+// The go-ethereum library is free software: you can redistribute it and/or modify
+// it under the terms of the GNU Lesser General Public License as published by
+// the Free Software Foundation, either version 3 of the License, or
+// (at your option) any later version.
+//
+// The go-ethereum library is distributed in the hope that it will be useful,
+// but WITHOUT ANY WARRANTY; without even the implied warranty of
+// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
+// GNU Lesser General Public License for more details.
+//
+// You should have received a copy of the GNU Lesser General Public License
+// along with the go-ethereum library. If not, see .
+
+package consensus
+
+import (
+ "bytes"
+ "math/big"
+ "strings"
+ "testing"
+
+ "github.com/XinFinOrg/XDPoSChain/common"
+ "github.com/XinFinOrg/XDPoSChain/core/types"
+ "github.com/XinFinOrg/XDPoSChain/log"
+)
+
+// stubCanonicalChain is the minimal chain the judgment needs: the canonical
+// header per height and which hashes have a body stored behind them. It is
+// both a CanonicalChain and a BlockStorer — the judgement must never look at
+// anything else.
+type stubCanonicalChain struct {
+ headers map[uint64]*types.Header
+ bodies map[common.Hash]bool
+}
+
+func (s *stubCanonicalChain) GetHeaderByNumber(number uint64) *types.Header {
+ return s.headers[number]
+}
+
+func (s *stubCanonicalChain) HasBlock(hash common.Hash, number uint64) bool {
+ return s.bodies[hash]
+}
+
+// stubHeaderOnlyChain is a CanonicalChain that is deliberately not a
+// BlockStorer, the way a header-only chain (e.g. *core.HeaderChain) looks to
+// the judgment: it must be skipped as SkipUnjudgeable, not judged "not stored".
+type stubHeaderOnlyChain struct {
+ headers map[uint64]*types.Header
+}
+
+func (s *stubHeaderOnlyChain) GetHeaderByNumber(number uint64) *types.Header {
+ return s.headers[number]
+}
+
+func TestShouldHandleProposedBlock(t *testing.T) {
+ const height = uint64(906)
+ canonical := &types.Header{Number: big.NewInt(int64(height)), Coinbase: common.BytesToAddress([]byte{0x01})}
+ fork := &types.Header{Number: big.NewInt(int64(height)), Coinbase: common.BytesToAddress([]byte{0x02})}
+
+ for _, c := range []struct {
+ name string
+ headers map[uint64]*types.Header
+ bodies map[common.Hash]bool
+ header *types.Header
+ storerless bool
+ wantOK bool
+ wantReason SkipReason
+ wantCanonHash common.Hash
+ }{
+ {
+ name: "no canonical header at height",
+ headers: map[uint64]*types.Header{},
+ bodies: map[common.Hash]bool{},
+ header: canonical,
+ wantOK: false,
+ wantReason: SkipNoCanonicalHeader,
+ wantCanonHash: common.Hash{},
+ },
+ {
+ name: "non-canonical",
+ headers: map[uint64]*types.Header{
+ height: canonical,
+ },
+ bodies: map[common.Hash]bool{
+ fork.Hash(): true,
+ },
+ header: fork,
+ wantOK: false,
+ wantReason: SkipNonCanonical,
+ wantCanonHash: canonical.Hash(),
+ },
+ {
+ name: "canonical header without stored body",
+ headers: map[uint64]*types.Header{
+ height: canonical,
+ },
+ bodies: map[common.Hash]bool{},
+ header: canonical,
+ wantOK: false,
+ wantReason: SkipBodyNotStored,
+ wantCanonHash: canonical.Hash(),
+ },
+ {
+ name: "canonical header with stored body",
+ headers: map[uint64]*types.Header{
+ height: canonical,
+ },
+ bodies: map[common.Hash]bool{
+ canonical.Hash(): true,
+ },
+ header: canonical,
+ wantOK: true,
+ wantReason: "",
+ wantCanonHash: canonical.Hash(),
+ },
+ {
+ name: "header-only chain cannot answer the storage half",
+ headers: map[uint64]*types.Header{
+ height: canonical,
+ },
+ header: canonical,
+ storerless: true,
+ wantOK: false,
+ wantReason: SkipUnjudgeable,
+ },
+ {
+ name: "header-only chain, absent height still unjudgeable",
+ headers: map[uint64]*types.Header{},
+ header: canonical,
+ storerless: true,
+ wantOK: false,
+ wantReason: SkipUnjudgeable,
+ },
+ } {
+ t.Run(c.name, func(t *testing.T) {
+ var chain CanonicalChain = &stubCanonicalChain{headers: c.headers, bodies: c.bodies}
+ if c.storerless {
+ chain = &stubHeaderOnlyChain{headers: c.headers}
+ }
+ ok, reason, canonicalHash := ShouldHandleProposedBlock(chain, c.header)
+ if ok != c.wantOK {
+ t.Errorf("ok = %v, want %v", ok, c.wantOK)
+ }
+ if reason != c.wantReason {
+ t.Errorf("reason = %q, want %q", reason, c.wantReason)
+ }
+ if canonicalHash != c.wantCanonHash {
+ t.Errorf("canonicalHash = %v, want %v", canonicalHash, c.wantCanonHash)
+ }
+ })
+ }
+}
+
+// TestSkipLogLevelGradesByReason pins the skip-log grading contract at
+// SkipLogLevel's definition site: SkipUnjudgeable — the wiring-bug skip that
+// halts QC processing and voting — surfaces as Error and sits outside the
+// Info/Warn discipline, the reorg-race skips (SkipNonCanonical,
+// SkipNoCanonicalHeader) surface as Warn, while SkipBodyNotStored — the one
+// routine sync skip — stays at Info, both here at the unit level and
+// end-to-end at the handler level (TestProposedBlockHandlerGradesSkipLogLevelByReason).
+// The empty reason means "accepted" and must not inflate the Warn level.
+// Any other unregistered reason grades Warn: a future skip reason whose
+// registration is forgotten must surface loudly, not hide at Info.
+
+func TestSkipLogLevelGradesByReason(t *testing.T) {
+ var logBuf bytes.Buffer
+ glog := log.NewGlogHandler(log.NewTerminalHandlerWithLevel(&logBuf, log.LevelInfo, false))
+ glog.Verbosity(log.LevelInfo)
+ prevLog := log.Root()
+ log.SetDefault(log.NewLogger(glog))
+ defer log.SetDefault(prevLog)
+
+ for _, c := range []struct {
+ reason SkipReason
+ want string
+ }{
+ {
+ SkipUnjudgeable,
+ "ERROR",
+ },
+ {SkipNonCanonical, "WARN"},
+ {SkipNoCanonicalHeader, "WARN"},
+ {SkipBodyNotStored, "INFO"},
+ {SkipReason(""), "INFO"},
+ // An unregistered reason must grade as the anomaly it most likely is,
+ // so a future reason whose registration is forgotten surfaces loudly
+ // instead of hiding at Info. The default branch is Warn for exactly
+ // this case.
+ {SkipReason("unregistered-future-reason"), "WARN"},
+ } {
+ logBuf.Reset()
+ SkipLogLevel(c.reason)("skip level probe")
+ level := ""
+ for _, line := range strings.Split(logBuf.String(), "\n") {
+ if strings.Contains(line, "skip level probe") {
+ level = strings.Fields(line)[0]
+ // glog pads the level to five characters, so "ERROR" carries
+ // no trailing space and runs straight into the timestamp.
+ if i := strings.IndexByte(level, '['); i >= 0 {
+ level = level[:i]
+ }
+ break
+ }
+ }
+ if level != c.want {
+ t.Errorf("reason %q graded to %q, want %s (log: %q)", c.reason, level, c.want, logBuf.String())
+ }
+ }
+}
diff --git a/consensus/tests/engine_v2_tests/proposed_block_test.go b/consensus/tests/engine_v2_tests/proposed_block_test.go
index 5f29b609ec02..badd21acc81b 100644
--- a/consensus/tests/engine_v2_tests/proposed_block_test.go
+++ b/consensus/tests/engine_v2_tests/proposed_block_test.go
@@ -1,14 +1,20 @@
package engine_v2_tests
import (
+ "bytes"
"fmt"
+ "math/big"
+ "strings"
"testing"
"time"
"github.com/XinFinOrg/XDPoSChain/accounts/abi/bind/backends"
+ "github.com/XinFinOrg/XDPoSChain/common"
+ "github.com/XinFinOrg/XDPoSChain/consensus"
"github.com/XinFinOrg/XDPoSChain/consensus/XDPoS"
"github.com/XinFinOrg/XDPoSChain/consensus/XDPoS/utils"
"github.com/XinFinOrg/XDPoSChain/core/types"
+ "github.com/XinFinOrg/XDPoSChain/log"
"github.com/XinFinOrg/XDPoSChain/params"
"github.com/stretchr/testify/assert"
)
@@ -318,18 +324,128 @@ func TestShouldNotSendVoteMsgIfBlockNotExtendedFromAncestor(t *testing.T) {
t.Fatal("Fail to decode extra data", err)
}
assert.Equal(t, types.Round(9), extraField.Round)
- // Set the lockQC and other pre-requist properties by block 906
+
+ // Process the QC carried by block 906 to set the lockQC without voting,
+ // so the negative case below is not blocked by the highestVotedRound
+ // voting-rule gate. The lockQC is the QC embedded in the block the QC
+ // points at, i.e. block 905's own QC pointing at block 904.
+ var extra906 types.ExtraFields_v2
+ err = utils.DecodeBytesExtraFields(currentBlock.Extra(), &extra906)
+ if err != nil {
+ t.Fatal("Fail to decode extra data of block 906", err)
+ }
+ err = engineV2.ProcessQCFaker(blockchain, extra906.QuorumCert)
+ if err != nil {
+ t.Fatal("Fail to process QC of block 906", err)
+ }
+
+ // Negative case: propose the canonical block 903, whose height is below
+ // the locked ancestor block 904. verifyVotingRule must reject it as not
+ // extending from the lockQC ancestor, so no vote is broadcast. The
+ // block is canonical on purpose, so the entry canonicality re-check of
+ // ProposedBlockHandler does not short-circuit the branch under test.
+ olderCanonicalBlock := blockchain.GetBlockByNumber(903)
+ assert.Equal(t, olderCanonicalBlock.Hash(), blockchain.GetCanonicalHash(olderCanonicalBlock.NumberU64()))
+ engineV2.SetNewRoundFaker(blockchain, types.Round(3), false)
+ err = engineV2.ProposedBlockHandler(blockchain, olderCanonicalBlock.Header())
+ if err != nil {
+ t.Fatal("Fail propose proposedBlock handler", err)
+ }
+ // Should not receive anything from the channel
+ select {
+ case <-engineV2.BroadcastCh:
+ t.Fatal("Should not trigger vote")
+ case <-time.After(3 * time.Second):
+ // Shoud not trigger setNewRound
+ round, _, _, _, _, _ := engineV2.GetPropertiesFaker()
+ assert.Equal(t, types.Round(3), round)
+ }
+
+ // Positive control: the canonical block 906 extends the locked ancestor,
+ // so its vote is broadcast as usual.
err = engineV2.ProposedBlockHandler(blockchain, currentBlock.Header())
if err != nil {
t.Fatal("Error while handling block 16", err)
}
vote := <-engineV2.BroadcastCh
assert.Equal(t, types.Round(6), vote.(*types.Vote).ProposedBlockInfo.Round)
+}
+
+/*
+ Block and round relationship diagram for this test
+ ... - 904(4) - 905(5) - 906(6) (canonical)
+ \ 904'(7) - 905'(8) - 906'(9) (fork)
+
+Unlike TestShouldNotSendVoteMsgIfBlockNotExtendedFromAncestor, whose
+proposed block is below the locked ancestor (blockNumDiff = -1, so the
+parent walk of isExtendingFromAncestor runs zero iterations), this test
+forces a positive height difference: the proposed block is canonical
+and higher than the locked ancestor, but the locked ancestor sits on a
+fork that was reorged away. The parent walk therefore runs and only the
+final hash comparison can reject the block.
+
+Note on reachability: only the negative branch of the parent walk is
+covered here. The positive branch (walk lands on the ancestor) needs a
+proposed block on the fork chain itself, which the entry canonicality
+re-check of ProposedBlockHandler now short-circuits, so it cannot be
+reached through the handler anymore.
+*/
+func TestShouldNotSendVoteMsgIfCanonicalBlockNotExtendedFromForkedAncestor(t *testing.T) {
+ skipLongInShortMode(t)
+ // Block number 905, 906 have forks and forkedBlock is the 906th
+ var numOfForks = new(int)
+ *numOfForks = 3
+ blockchain, _, currentBlock, _, _, forkedBlock := PrepareXDCTestBlockChainForV2Engine(t, 906, params.TestXDPoSMockChainConfig, &ForkedBlockOptions{numOfForkedBlocks: numOfForks})
+ engineV2 := blockchain.Engine().(*XDPoS.XDPoS).EngineV2
+
+ forkedAncestor := blockchain.GetBlockByHash(blockchain.GetBlockByHash(forkedBlock.ParentHash()).ParentHash())
+ assert.NotEqual(t, blockchain.GetCanonicalHash(forkedAncestor.NumberU64()), forkedAncestor.Hash())
+
+ // Process the QC carried by the canonical block 906, then the QC carried
+ // by the forked block 906', to set the lockQC without voting, so the
+ // negative case below is not blocked by the highestVotedRound voting-rule
+ // gate. The lockQC is the QC embedded in the block the processed QC
+ // points at: the first call leaves it at the canonical block 905's own
+ // QC (pointing at the canonical block 904), and the second call replaces
+ // it with the forked block 905's own QC, which points at the forked
+ // block 904'.
+ var extra906 types.ExtraFields_v2
+ err := utils.DecodeBytesExtraFields(currentBlock.Extra(), &extra906)
+ if err != nil {
+ t.Fatal("Fail to decode extra data of block 906", err)
+ }
+ err = engineV2.ProcessQCFaker(blockchain, extra906.QuorumCert)
+ if err != nil {
+ t.Fatal("Fail to process QC of block 906", err)
+ }
+
+ var extraForked906 types.ExtraFields_v2
+ err = utils.DecodeBytesExtraFields(forkedBlock.Extra(), &extraForked906)
+ if err != nil {
+ t.Fatal("Fail to decode extra data of forked block 906'", err)
+ }
+ err = engineV2.ProcessQCFaker(blockchain, extraForked906.QuorumCert)
+ if err != nil {
+ t.Fatal("Fail to process QC of forked block 906'", err)
+ }
- // Find the first forked block at block 14th
- firstForkedBlock := blockchain.GetBlockByHash(blockchain.GetBlockByHash(forkedBlock.ParentHash()).ParentHash())
- engineV2.SetNewRoundFaker(blockchain, types.Round(7), false)
- err = engineV2.ProposedBlockHandler(blockchain, firstForkedBlock.Header())
+ // Pin the preconditions the negative case below relies on: the lockQC
+ // points at the forked ancestor and the current round leaves room for a
+ // vote on the canonical block 906 (round 6).
+ _, lockQC, _, _, highestVotedRound, _ := engineV2.GetPropertiesFaker()
+ if assert.NotNil(t, lockQC) {
+ assert.Equal(t, forkedAncestor.Hash(), lockQC.ProposedBlockInfo.Hash)
+ }
+ assert.Equal(t, types.Round(0), highestVotedRound)
+
+ // Negative case: propose the canonical block 906. It passes the entry
+ // canonicality re-check of ProposedBlockHandler on purpose, so the branch
+ // under test is verifyVotingRule: the block's QC round does not outrank
+ // the lockQC round, so isExtendingFromAncestor walks two parents down to
+ // the canonical block 904, which is not the forked ancestor 904', and the
+ // vote must be dropped.
+ engineV2.SetNewRoundFaker(blockchain, types.Round(6), false)
+ err = engineV2.ProposedBlockHandler(blockchain, currentBlock.Header())
if err != nil {
t.Fatal("Fail propose proposedBlock handler", err)
}
@@ -338,9 +454,12 @@ func TestShouldNotSendVoteMsgIfBlockNotExtendedFromAncestor(t *testing.T) {
case <-engineV2.BroadcastCh:
t.Fatal("Should not trigger vote")
case <-time.After(3 * time.Second):
- // Shoud not trigger setNewRound
- round, _, _, _, _, _ := engineV2.GetPropertiesFaker()
- assert.Equal(t, types.Round(7), round)
+ // Should not trigger setNewRound
+ round, lockQC, _, _, _, _ := engineV2.GetPropertiesFaker()
+ assert.Equal(t, types.Round(6), round)
+ if assert.NotNil(t, lockQC) {
+ assert.Equal(t, forkedAncestor.Hash(), lockQC.ProposedBlockInfo.Hash)
+ }
}
}
@@ -396,3 +515,366 @@ func TestProposedBlockMessageHandlerNotGenerateVoteIfSignerNotInMNlist(t *testin
assert.Equal(t, types.Round(6), round)
}
}
+
+// TestProposedBlockHandlerSkipsNonCanonicalBlock pins the canonicality and
+// storage re-check directly in front of processQC: the callers' gates race
+// with concurrent imports, and processQC updates highestQuorumCert, the lock
+// QC and the commit block before its own existence check, which a stored
+// side-chain block passes. The forked block below is exactly what the
+// window leaves behind: stored, with a valid parent QC, but no longer
+// canonical at its height.
+func TestProposedBlockHandlerSkipsNonCanonicalBlock(t *testing.T) {
+ var numOfForks = new(int)
+ *numOfForks = 1
+ blockchain, _, currentBlock, _, _, forkedBlock := PrepareXDCTestBlockChainForV2Engine(t, 906, params.TestXDPoSMockChainConfig, &ForkedBlockOptions{numOfForkedBlocks: numOfForks})
+ engineV2 := blockchain.Engine().(*XDPoS.XDPoS).EngineV2
+
+ // Precondition: the fork is stored but is not canonical at its height.
+ assert.NotNil(t, blockchain.GetBlockByHash(forkedBlock.Hash()))
+ assert.NotEqual(t, forkedBlock.Hash(), blockchain.GetCanonicalHash(forkedBlock.NumberU64()))
+ assert.Equal(t, currentBlock.Hash(), blockchain.GetCanonicalHash(forkedBlock.NumberU64()))
+
+ beforeRound, beforeLockQC, beforeHighestQC, beforeTimeoutCert, beforeVotedRound, beforeCommit := engineV2.GetPropertiesFaker()
+
+ err := engineV2.ProposedBlockHandler(blockchain, forkedBlock.Header())
+ assert.Nil(t, err)
+
+ round, lockQC, highestQC, timeoutCert, votedRound, commit := engineV2.GetPropertiesFaker()
+ assert.Equal(t, beforeRound, round)
+ assert.Equal(t, beforeLockQC, lockQC)
+ assert.Equal(t, beforeHighestQC, highestQC)
+ assert.Equal(t, beforeTimeoutCert, timeoutCert)
+ assert.Equal(t, beforeVotedRound, votedRound)
+ assert.Equal(t, beforeCommit, commit)
+
+ // A non-canonical block must not reach the vote broadcast either.
+ select {
+ case vote := <-engineV2.BroadcastCh:
+ t.Fatalf("non-canonical block must not trigger a vote, got round %v", vote.(*types.Vote).ProposedBlockInfo.Round)
+ case <-time.After(2 * time.Second):
+ }
+}
+
+// reorgingChainReader simulates a concurrent reorg landing inside the
+// handler: x.lock does not block InsertChain, so the canonical answer for
+// a height can change between two reads. Reads of other heights, and the
+// first truthfulReads reads of the watched height, are answered by the
+// wrapped chain; later reads of the watched height return the fork header
+// that took the height over.
+type reorgingChainReader struct {
+ consensus.ChainReader
+ watchNumber uint64
+ truthfulReads int
+ forkHeader *types.Header
+
+ // reads counts the reads of the watched height so the injection knows
+ // when to switch to the fork header; it is implementation detail of the
+ // wrapper, not an assertion anchor.
+ reads int
+}
+
+func (r *reorgingChainReader) GetHeaderByNumber(number uint64) *types.Header {
+ if number == r.watchNumber {
+ r.reads++
+ if r.reads > r.truthfulReads {
+ return r.forkHeader
+ }
+ }
+ return r.ChainReader.GetHeaderByNumber(number)
+}
+
+// HasBlock forwards the storage half of the proposed-block judgment to the
+// wrapped chain, which must be a consensus.BlockStorer — the wrapper only
+// rewrites canonicality, never storage.
+func (r *reorgingChainReader) HasBlock(hash common.Hash, number uint64) bool {
+ return r.ChainReader.(consensus.BlockStorer).HasBlock(hash, number)
+}
+
+// TestProposedBlockHandlerSkipsReorgedBlockBeforeProcessQC covers the first
+// re-check point, directly in front of processQC: x.lock does not block
+// InsertChain, so a concurrent import can take the height over before the
+// handler reaches it. The wrapper serves the fork header from the first read
+// of the watched height on, making the reorg deterministic instead of racing
+// a real InsertChain. Block 906's embedded QC certifies block 905 at round 5,
+// higher than the engine's initial highestQuorumCert (round 0), so any
+// processQC run would visibly move the engine state — exactly what the gate
+// must prevent for a reorged-away block.
+func TestProposedBlockHandlerSkipsReorgedBlockBeforeProcessQC(t *testing.T) {
+ skipLongInShortMode(t)
+ var numOfForks = new(int)
+ *numOfForks = 1
+ blockchain, _, currentBlock, _, _, forkedBlock := PrepareXDCTestBlockChainForV2Engine(t, 906, params.TestXDPoSMockChainConfig, &ForkedBlockOptions{numOfForkedBlocks: numOfForks})
+ engineV2 := blockchain.Engine().(*XDPoS.XDPoS).EngineV2
+
+ // Precondition: the fork sits at the same height but is not canonical.
+ assert.NotNil(t, blockchain.GetBlockByHash(forkedBlock.Hash()))
+ assert.Equal(t, forkedBlock.NumberU64(), currentBlock.NumberU64())
+ assert.NotEqual(t, forkedBlock.Hash(), blockchain.GetCanonicalHash(forkedBlock.NumberU64()))
+
+ // The height is already reorged when the handler first reads it: zero
+ // truthful reads.
+ chain := &reorgingChainReader{
+ ChainReader: blockchain,
+ watchNumber: currentBlock.NumberU64(),
+ truthfulReads: 0,
+ forkHeader: forkedBlock.Header(),
+ }
+ beforeRound, beforeLockQC, beforeHighQC, beforeHighTC, beforeVotedRound, beforeCommitBlock := engineV2.GetPropertiesFaker()
+ err := engineV2.ProposedBlockHandler(chain, currentBlock.Header())
+ assert.Nil(t, err)
+
+ // The block was reorged away before processQC: the engine state must be
+ // untouched — no QC processed, nothing locked, nothing committed, no
+ // round advanced.
+ afterRound, afterLockQC, afterHighQC, afterHighTC, afterVotedRound, afterCommitBlock := engineV2.GetPropertiesFaker()
+ assert.Equal(t, beforeRound, afterRound)
+ assert.Equal(t, beforeLockQC, afterLockQC)
+ assert.Equal(t, beforeHighQC, afterHighQC)
+ assert.Equal(t, beforeHighTC, afterHighTC)
+ assert.Equal(t, beforeVotedRound, afterVotedRound)
+ assert.Equal(t, beforeCommitBlock, afterCommitBlock)
+
+ // Nothing may be broadcast either.
+ select {
+ case vote := <-engineV2.BroadcastCh:
+ t.Fatalf("a block reorged away before processQC must not be voted for, got vote for round %v hash %v", vote.(*types.Vote).ProposedBlockInfo.Round, vote.(*types.Vote).ProposedBlockInfo.Hash)
+ case <-time.After(2 * time.Second):
+ }
+
+ // No read-count anchor: truthfulReads is 0, so every read of the watched
+ // height sees the fork and the handler skips before processQC no matter
+ // how many reads are added or removed around the re-check — the state
+ // anchors above cover the drift. If this test fails after touching the
+ // handler's chain reads, re-count which read consumes the fork header
+ // and move the injection (truthfulReads), do not just bump the number.
+}
+
+// TestProposedBlockHandlerDropsVoteForReorgedBlock covers the second
+// re-check point, right before sendVote: x.lock does not block InsertChain,
+// so a concurrent import can take the height over after processQC ran on the
+// still canonical block but before the vote is broadcast. The wrapper serves
+// the real canonical header to the pre-processQC re-check and the fork
+// header from then on, making the mid-handler reorg deterministic instead of
+// racing a real InsertChain. processQC legitimately runs in this window (the
+// block was canonical when it read the chain, so its state write is
+// expected); the vote is what must be dropped.
+func TestProposedBlockHandlerDropsVoteForReorgedBlock(t *testing.T) {
+ skipLongInShortMode(t)
+ var numOfForks = new(int)
+ *numOfForks = 1
+ // Height 906, not 901: the state anchor below needs a processQC run
+ // that visibly moves the engine. Block 901's embedded QC certifies
+ // block 900 at round 0, which the engine already holds — processQC
+ // would be a complete state no-op there. Block 906's QC certifies
+ // block 905 at round 5, strictly above the engine's initial state,
+ // so a successful processQC run advances highestQuorumCert, lockQC,
+ // currentRound and the commit block and the anchor can detect it.
+ blockchain, _, currentBlock, _, _, forkedBlock := PrepareXDCTestBlockChainForV2Engine(t, 906, params.TestXDPoSMockChainConfig, &ForkedBlockOptions{numOfForkedBlocks: numOfForks})
+ engineV2 := blockchain.Engine().(*XDPoS.XDPoS).EngineV2
+
+ // Precondition: the fork sits at the same height but is not canonical.
+ assert.NotNil(t, blockchain.GetBlockByHash(forkedBlock.Hash()))
+ assert.Equal(t, forkedBlock.NumberU64(), currentBlock.NumberU64())
+ assert.NotEqual(t, forkedBlock.Hash(), blockchain.GetCanonicalHash(forkedBlock.NumberU64()))
+
+ // One truthful read for the pre-processQC re-check; every later read of
+ // this height sees the reorg.
+ chain := &reorgingChainReader{
+ ChainReader: blockchain,
+ watchNumber: currentBlock.NumberU64(),
+ truthfulReads: 1,
+ forkHeader: forkedBlock.Header(),
+ }
+ beforeRound, _, beforeHighQC, _, _, _ := engineV2.GetPropertiesFaker()
+ err := engineV2.ProposedBlockHandler(chain, currentBlock.Header())
+ assert.Nil(t, err)
+
+ // State anchor: processQC must have run and moved the engine between
+ // the two gates. With truthfulReads 1 the pre-processQC re-check still
+ // sees the canonical header, so a read added in front of that re-check
+ // consumes the fork header there and skips the handler before
+ // processQC — the engine state then stays untouched and the assertions
+ // below fail loudly. Only a genuine pre-vote reorg lets processQC
+ // advance the state while the vote is still dropped.
+ afterRound, afterLockQC, afterHighQC, _, _, afterCommitBlock := engineV2.GetPropertiesFaker()
+ assert.Greater(t, afterRound, beforeRound, "processQC must have advanced the current round")
+ assert.Greater(t, afterHighQC.ProposedBlockInfo.Round, beforeHighQC.ProposedBlockInfo.Round, "processQC must have advanced highestQuorumCert")
+ assert.NotNil(t, afterLockQC, "processQC must have locked a parent QC")
+ assert.NotNil(t, afterCommitBlock, "processQC must have committed a block")
+
+ // The block was reorged away before the vote: nothing may be broadcast.
+ select {
+ case vote := <-engineV2.BroadcastCh:
+ t.Fatalf("a block reorged away before the vote must not be voted for, got vote for round %v hash %v", vote.(*types.Vote).ProposedBlockInfo.Round, vote.(*types.Vote).ProposedBlockInfo.Hash)
+ case <-time.After(2 * time.Second):
+ }
+
+}
+
+// TestProposedBlockHandlerSkipsBlockWithoutBody covers the storage half of
+// the gate: canonicality is a property of the header, not of the block, so
+// the fast sync header phase marks a height canonical before its body lands,
+// and the fetcher calls this handler after an insertBlock that can return
+// nil without writing anything (fast sync, and the downloadingBlock
+// short circuit). Only the downloader gates on storage, so the other
+// callers are covered here: a canonical header without a body must not
+// reach processQC or the vote broadcast, and must leave the engine state
+// untouched.
+func TestProposedBlockHandlerSkipsBlockWithoutBody(t *testing.T) {
+ skipLongInShortMode(t)
+ blockchain, _, currentBlock, signer, signFn, _ := PrepareXDCTestBlockChainForV2Engine(t, 906, params.TestXDPoSMockChainConfig, nil)
+ engineV2 := blockchain.Engine().(*XDPoS.XDPoS).EngineV2
+
+ // Insert a valid header for height 907 header-only, the way the fast
+ // sync header phase does: the height becomes canonical before any body
+ // is written. The block is built with the chain's own config so the
+ // header passes ValidateHeaderChain.
+ testConfig := legacyExecutionConfigForV2Tests(params.TestXDPoSMockChainConfig)
+ header907 := CreateBlock(blockchain, testConfig, currentBlock, 907, 7, signer.Hex(), signer, signFn, nil, nil, "").Header()
+ if _, err := blockchain.InsertHeaderChain([]*types.Header{header907}, 0); err != nil {
+ t.Fatal("Fail to insert header chain", err)
+ }
+ assert.Equal(t, header907.Hash(), blockchain.GetCanonicalHash(907))
+ assert.Nil(t, blockchain.GetBlock(header907.Hash(), 907))
+
+ before, beforeLockQC, beforeHighQC, beforeHighTC, beforeVotedRound, beforeCommitBlock := engineV2.GetPropertiesFaker()
+ err := engineV2.ProposedBlockHandler(blockchain, header907)
+ if err != nil {
+ t.Fatal("Fail propose proposedBlock handler", err)
+ }
+ // Should not receive anything from the channel
+ select {
+ case <-engineV2.BroadcastCh:
+ t.Fatal("Should not trigger vote")
+ case <-time.After(3 * time.Second):
+ }
+
+ // The engine state must be untouched: no QC processed, nothing voted,
+ // no round advanced.
+ after, afterLockQC, afterHighQC, afterHighTC, afterVotedRound, afterCommitBlock := engineV2.GetPropertiesFaker()
+ assert.Equal(t, before, after)
+ assert.Equal(t, beforeLockQC, afterLockQC)
+ assert.Equal(t, beforeHighQC, afterHighQC)
+ assert.Equal(t, beforeHighTC, afterHighTC)
+ assert.Equal(t, beforeVotedRound, afterVotedRound)
+ assert.Equal(t, beforeCommitBlock, afterCommitBlock)
+}
+
+// TestProposedBlockHandlerGradesSkipLogLevelByReason fixates the skip-log
+// grading: the reorg-race skips (SkipNonCanonical, SkipNoCanonicalHeader)
+// surface as Warn — the fast sync header phase marks heights canonical, so
+// a missing marker is never a routine sync state — while SkipBodyNotStored,
+// the one routine sync skip, stays at Info so a Warn per header-only fast
+// sync height cannot drown the level reserved for anomalies.
+func TestProposedBlockHandlerGradesSkipLogLevelByReason(t *testing.T) {
+ var numOfForks = new(int)
+ *numOfForks = 1
+ blockchain, _, currentBlock, signer, signFn, forkedBlock := PrepareXDCTestBlockChainForV2Engine(t, 906, params.TestXDPoSMockChainConfig, &ForkedBlockOptions{numOfForkedBlocks: numOfForks})
+ engineV2 := blockchain.Engine().(*XDPoS.XDPoS).EngineV2
+
+ // Capture the handler's own skip logs and assert on the level prefix of
+ // the specific skip line, so unrelated background logs cannot interfere.
+ var logBuf bytes.Buffer
+ prevLog := log.Root()
+ glog := log.NewGlogHandler(log.NewTerminalHandlerWithLevel(&logBuf, log.LevelInfo, false))
+ glog.Verbosity(log.LevelInfo)
+ log.SetDefault(log.NewLogger(glog))
+ defer log.SetDefault(prevLog)
+
+ levelOf := func(msg string) string {
+ for _, line := range strings.Split(logBuf.String(), "\n") {
+ if strings.Contains(line, msg) {
+ return strings.Fields(line)[0]
+ }
+ }
+ return ""
+ }
+
+ // A fork block is a genuine reorg-race skip: it must stay at Warn.
+ err := engineV2.ProposedBlockHandler(blockchain, forkedBlock.Header())
+ assert.Nil(t, err)
+ assert.Equal(t, "WARN", levelOf("skip block before processQC"),
+ "a non-canonical skip is a reorg race and must stay at Warn, have %q", logBuf.String())
+
+ // A height with no canonical header at all is a reorg-race skip too:
+ // markers only go missing above a contested head (a fork growing past
+ // the local chain, a displaced tip re-delivered, or a concurrent reorg),
+ // so it must stay at Warn. Reuse the fork header's extra data (so
+ // getExtraFields still parses) at a height nothing occupies.
+ noCanonical := *forkedBlock.Header()
+ noCanonical.Number = big.NewInt(99999)
+ logBuf.Reset()
+ err = engineV2.ProposedBlockHandler(blockchain, &noCanonical)
+ assert.Nil(t, err)
+ assert.Equal(t, "WARN", levelOf("skip block before processQC"),
+ "a no-canonical-header skip is a reorg race and must stay at Warn, have %q", logBuf.String())
+
+ // A canonical header without a stored body is the one routine skip: the
+ // fast sync header phase marks a height canonical before its body lands,
+ // so it must stay at Info. Build the shape the way that phase does —
+ // InsertHeaderChain one height past the canonical tip — mirroring
+ // TestProposedBlockHandlerSkipsBlockWithoutBody, and assert the level of
+ // the very same skip line the two reorg-race sections above graded to Warn.
+ testConfig := legacyExecutionConfigForV2Tests(params.TestXDPoSMockChainConfig)
+ header907 := CreateBlock(blockchain, testConfig, currentBlock, 907, 7, signer.Hex(), signer, signFn, nil, nil, "").Header()
+ if _, err := blockchain.InsertHeaderChain([]*types.Header{header907}, 0); err != nil {
+ t.Fatal("Fail to insert header chain", err)
+ }
+ assert.Equal(t, header907.Hash(), blockchain.GetCanonicalHash(907))
+ assert.Nil(t, blockchain.GetBlock(header907.Hash(), 907))
+ logBuf.Reset()
+ err = engineV2.ProposedBlockHandler(blockchain, header907)
+ assert.Nil(t, err)
+ assert.Equal(t, "INFO", levelOf("skip block before processQC"),
+ "a body-not-stored skip is routine fast-sync state and must stay at Info, have %q", logBuf.String())
+}
+
+// chainReaderOnly hides the BlockStorer half of a full chain: it is the
+// shape a wrapped or replaced chain type presents to the handler once the
+// storage half drops off the type. consensus.ShouldHandleProposedBlock must
+// skip it as SkipUnjudgeable at any height.
+type chainReaderOnly struct {
+ consensus.ChainReader
+}
+
+// TestProposedBlockHandlerSkipsUnjudgeableChain pins the "cannot judge is a
+// skip, not a failure" contract: a chain whose type hides the BlockStorer
+// half skips as SkipUnjudgeable — graded Error, because it means the node
+// stops processing QCs and voting — and returns nil. The fetcher's import
+// loop treats any handler error as an import failure and suppresses the
+// broadcast of an already imported block, so the skip must never surface as
+// an error.
+func TestProposedBlockHandlerSkipsUnjudgeableChain(t *testing.T) {
+ blockchain, _, currentBlock, _, _, _ := PrepareXDCTestBlockChainForV2Engine(t, 901, params.TestXDPoSMockChainConfig, nil)
+ engineV2 := blockchain.Engine().(*XDPoS.XDPoS).EngineV2
+
+ // Capture the handler's logs and assert the unjudgeable skip surfaces
+ // at Error — it is a liveness halt, not an ordinary reorg-race skip.
+ var logBuf bytes.Buffer
+ prevLog := log.Root()
+ glog := log.NewGlogHandler(log.NewTerminalHandlerWithLevel(&logBuf, log.LevelInfo, false))
+ glog.Verbosity(log.LevelInfo)
+ log.SetDefault(log.NewLogger(glog))
+ defer log.SetDefault(prevLog)
+ err := engineV2.ProposedBlockHandler(chainReaderOnly{blockchain}, currentBlock.Header())
+ assert.Nil(t, err, "an unjudgeable chain must not surface an error to the caller")
+
+ // The skip must surface at Error — a liveness halt, not a routine skip —
+ // and must not fall through to processQC or the vote path.
+ found := false
+ for _, line := range strings.Split(logBuf.String(), "\n") {
+ if strings.Contains(line, "skip block before processQC") {
+ assert.True(t, strings.HasPrefix(line, "ERROR"), "unjudgeable skip must log at Error, got line %q", line)
+ assert.Contains(t, line, "unjudgeable", "unjudgeable skip must carry the reason, got line %q", line)
+ found = true
+ break
+ }
+ }
+ assert.True(t, found, "unjudgeable skip line not found in logs: %q", logBuf.String())
+
+ select {
+ case vote := <-engineV2.BroadcastCh:
+ t.Fatalf("unjudgeable chain must not drive the vote path, got %v", vote)
+ default:
+ }
+}
diff --git a/core/blockchain.go b/core/blockchain.go
index 54f5291c8895..a4327f3aea4a 100644
--- a/core/blockchain.go
+++ b/core/blockchain.go
@@ -229,6 +229,15 @@ type BlockChain struct {
finalizedTrade *lru.Cache[common.Hash, interface{}] // include both trades which force update to closed/liquidated by the protocol
}
+// The full chain must be able to answer both halves of the proposed-block
+// judgment; a compile-time check keeps the storage half from silently
+// dropping off the type (see consensus.ShouldHandleProposedBlock, whose
+// CanonicalChain parameter alone would only catch it at runtime).
+var _ interface {
+ consensus.CanonicalChain
+ consensus.BlockStorer
+} = (*BlockChain)(nil)
+
type blockchainOpenConfig struct {
readOnly bool
chainConfig *params.ChainConfig
diff --git a/core/headerchain.go b/core/headerchain.go
index 2fa0a4d635d4..3a5b255c1145 100644
--- a/core/headerchain.go
+++ b/core/headerchain.go
@@ -551,7 +551,16 @@ func (hc *HeaderChain) Config() *params.ChainConfig { return hc.config }
func (hc *HeaderChain) Engine() consensus.Engine { return hc.engine }
// GetBlock implements consensus.ChainReader, and returns nil for every input as
-// a header chain does not have blocks available for retrieval.
+// a header chain does not have blocks available for retrieval. The stub exists
+// only to close the interface: ValidateHeaderChain is the only place a HeaderChain
+// is handed to the consensus engine as a ChainReader, and only a consensus
+// decision that leans on block storage (e.g. consensus.ShouldHandleProposedBlock)
+// would be hurt by such a stub — every block would fail its storage half. The
+// judgment therefore takes its storage input through the optional
+// consensus.BlockStorer interface, which a header chain deliberately does not
+// implement: handing a HeaderChain to the judgment skips every block as
+// SkipUnjudgeable — graded Error, so it fails loudly rather than a silent
+// "not stored" for every block.
func (hc *HeaderChain) GetBlock(hash common.Hash, number uint64) *types.Block {
return nil
}
diff --git a/core/headerchain_test.go b/core/headerchain_test.go
new file mode 100644
index 000000000000..bae4a71d6df8
--- /dev/null
+++ b/core/headerchain_test.go
@@ -0,0 +1,50 @@
+package core
+
+import (
+ "math/big"
+ "testing"
+
+ "github.com/XinFinOrg/XDPoSChain/common"
+ "github.com/XinFinOrg/XDPoSChain/consensus"
+ "github.com/XinFinOrg/XDPoSChain/core/rawdb"
+ "github.com/XinFinOrg/XDPoSChain/core/types"
+ "github.com/XinFinOrg/XDPoSChain/params"
+)
+
+// TestHeaderChainShouldHandleProposedBlock fixates the interface adaptation:
+// *HeaderChain satisfies consensus.ChainReader only so the interface stays
+// closed, but it stores no block bodies and deliberately does not implement
+// consensus.BlockStorer. The shared proposed-block gate must therefore skip
+// every block with SkipUnjudgeable — graded Error, not judge every block
+// "not stored" — so a HeaderChain handed to the consensus engine by mistake
+// cannot silently drop QC processing and voting behind an Info log.
+func TestHeaderChainShouldHandleProposedBlock(t *testing.T) {
+ db := rawdb.NewMemoryDatabase()
+ genesis := &types.Header{Number: big.NewInt(0)}
+ rawdb.WriteHeader(db, genesis)
+ rawdb.WriteCanonicalHash(db, genesis.Hash(), 0)
+ hc, err := NewHeaderChain(db, params.TestChainConfig, nil, func() bool { return false })
+ if err != nil {
+ t.Fatal("Fail to create header chain", err)
+ }
+
+ // A header stored header-only (the fast sync header phase shape): canonical
+ // at its height, so the gate reaches the storage half — which a header
+ // chain cannot answer — and must skip it as unjudgeable.
+ if ok, reason, _ := consensus.ShouldHandleProposedBlock(hc, genesis); ok || reason != consensus.SkipUnjudgeable {
+ t.Fatalf("a header chain must be skipped as unjudgeable by the proposed-block gate, got ok=%v reason=%q", ok, reason)
+ }
+
+ missing := &types.Header{Number: big.NewInt(5)}
+ // The BlockStorer assertion runs before any chain read, so an absent
+ // height is reported as unjudgeable too, not silently skipped on the
+ // canonicality half.
+ if ok, reason, canonicalHash := consensus.ShouldHandleProposedBlock(hc, missing); ok || reason != consensus.SkipUnjudgeable || canonicalHash != (common.Hash{}) {
+ t.Fatalf("an absent height must also be reported as unjudgeable on a header chain, got ok=%v reason=%q canonicalHash=%v", ok, reason, canonicalHash)
+ }
+
+ // The block-retrieval stub itself.
+ if block := hc.GetBlock(genesis.Hash(), 0); block != nil {
+ t.Fatalf("HeaderChain.GetBlock must return nil, got %v", block)
+ }
+}
diff --git a/eth/downloader/downloader.go b/eth/downloader/downloader.go
index 6e280b76eb1d..2c9b9cf086ba 100644
--- a/eth/downloader/downloader.go
+++ b/eth/downloader/downloader.go
@@ -27,6 +27,7 @@ import (
"github.com/XinFinOrg/XDPoSChain"
"github.com/XinFinOrg/XDPoSChain/common"
+ "github.com/XinFinOrg/XDPoSChain/consensus"
"github.com/XinFinOrg/XDPoSChain/consensus/XDPoS/engines/engine_v2"
"github.com/XinFinOrg/XDPoSChain/core/rawdb"
"github.com/XinFinOrg/XDPoSChain/core/state"
@@ -212,6 +213,9 @@ type BlockChain interface {
// InsertChain inserts a batch of blocks into the local chain.
InsertChain(types.Blocks) (int, error)
+ // GetHeaderByNumber retrieves a canonical header from the local chain by height.
+ GetHeaderByNumber(number uint64) *types.Header
+
// InterruptInsert disables or enables chain insertion.
InterruptInsert(on bool)
@@ -1661,11 +1665,34 @@ func (d *Downloader) importBlockResults(results []*fetchResult) error {
}
return fmt.Errorf("%w: %v", errInvalidChain, err)
}
+ // A nil InsertChain does not mean the tail reached the canonical chain:
+ // a fork batch is stored as side entries and a parked tail is not stored
+ // at all. This pre-filter is an early exit plus the observation point for
+ // routine skips: fork tails are expected during sync, so they are logged
+ // here at Info, while the handler's own skip log stays a Warn for rare
+ // reorg races. It runs the same judgment on the same chain as the
+ // handler's first check, so it cannot diverge from it; correctness still
+ // rests on the handler's re-checks, since the chain can move after this
+ // point. The level is therefore graded by call site, not by reason: the
+ // same SkipReason is Info here (routine sync tails) and Warn at the
+ // handler (rare races) by design — do not "fix" the divergence to
+ // consensus.SkipLogLevel, it would turn every sync fork tail into a Warn.
if d.handleProposedBlock != nil {
- header := blocks[len(blocks)-1].Header()
- err := d.handleProposedBlock(header)
- if err != nil {
- log.Info("[downloader] handle proposed block has error", "err", err, "block hash", header.Hash(), "number", header.Number)
+ tail := blocks[len(blocks)-1]
+ ok, reason, canonicalHash := consensus.ShouldHandleProposedBlock(d.blockchain, tail.Header())
+ if ok {
+ if err := d.handleProposedBlock(tail.Header()); err != nil {
+ log.Info("[downloader] handle proposed block has error", "err", err, "block hash", tail.Hash(), "number", tail.Number())
+ }
+ } else if reason == consensus.SkipUnjudgeable {
+ // The chain cannot answer the storage half of the judgment; failing
+ // the sync over a wiring bug would be worse than skipping the
+ // handler, whose correctness still rests on the engine's own
+ // re-checks. Like the skips below, this is graded by call site:
+ // an explicit Error, not consensus.SkipLogLevel.
+ log.Error("[downloader] cannot judge proposed block", "reason", reason, "block hash", tail.Hash(), "number", tail.Number())
+ } else {
+ log.Info("[downloader] skipped proposed block handler", "block hash", tail.Hash(), "number", tail.Number(), "reason", reason, "canonicalHash", canonicalHash)
}
}
return nil
diff --git a/eth/downloader/downloader_test.go b/eth/downloader/downloader_test.go
index 02f4ab9b3733..fb220de10453 100644
--- a/eth/downloader/downloader_test.go
+++ b/eth/downloader/downloader_test.go
@@ -68,13 +68,28 @@ type downloadTester struct {
peerDb ethdb.Database // Database of the peers containing all data
peers map[string]*downloadTesterPeer
- ownHashes []common.Hash // Hash chain belonging to the tester
- ownHeaders map[common.Hash]*types.Header // Headers belonging to the tester
- ownBlocks map[common.Hash]*types.Block // Blocks belonging to the tester
- ownReceipts map[common.Hash]types.Receipts // Receipts belonging to the tester
- ownChainTd map[common.Hash]*big.Int // Total difficulties of the blocks in the local chain
+ ownHashes []common.Hash // Canonical hash chain in height order, rebuilt from ownCanonical
+ ownHeaders map[common.Hash]*types.Header // Headers belonging to the tester
+ ownBlocks map[common.Hash]*types.Block // Blocks belonging to the tester
+ ownReceipts map[common.Hash]types.Receipts // Receipts belonging to the tester
+ ownChainTd map[common.Hash]*big.Int // Total difficulties of the blocks in the local chain
+ ownCanonical map[uint64]common.Hash // Canonical number-to-hash table, mirroring the real chain's marker writes
insertHeaderChainHook func([]*types.Header) error
+ // parkTailOnce makes the next InsertChain park the batch tail, the way
+ // insertChain parks future blocks: everything but the tail is written.
+ // Consumed by that one call.
+ parkTailOnce bool
+ // extendTailAfterInsert, when set, is called after the next InsertChain
+ // with the batch tail and stores the child it returns, simulating a
+ // concurrent import landing before importBlockResults reads the head.
+ // Consumed by that one call.
+ extendTailAfterInsert func(tail *types.Block) *types.Block
+
+ // proposedCalls and lastProposedHeader record handleProposedBlock calls
+ // for assertions; read them together via proposedState.
+ proposedCalls int
+ lastProposedHeader *types.Header
// headHeaderCap, when non-zero, caps the height reported by CurrentHeader.
// It models the real chain, where importing blocks moves the header head
@@ -104,14 +119,15 @@ func newTester() *downloadTester {
// shared test genesis.
func newTesterWithGenesis(genesis *types.Block, peerDb ethdb.Database) *downloadTester {
tester := &downloadTester{
- genesis: genesis,
- peerDb: peerDb,
- peers: make(map[string]*downloadTesterPeer),
- ownHashes: []common.Hash{genesis.Hash()},
- ownHeaders: map[common.Hash]*types.Header{genesis.Hash(): genesis.Header()},
- ownBlocks: map[common.Hash]*types.Block{genesis.Hash(): genesis},
- ownReceipts: map[common.Hash]types.Receipts{genesis.Hash(): nil},
- ownChainTd: map[common.Hash]*big.Int{genesis.Hash(): genesis.Difficulty()},
+ genesis: genesis,
+ peerDb: peerDb,
+ peers: make(map[string]*downloadTesterPeer),
+ ownHashes: []common.Hash{genesis.Hash()},
+ ownHeaders: map[common.Hash]*types.Header{genesis.Hash(): genesis.Header()},
+ ownBlocks: map[common.Hash]*types.Block{genesis.Hash(): genesis},
+ ownReceipts: map[common.Hash]types.Receipts{genesis.Hash(): nil},
+ ownChainTd: map[common.Hash]*big.Int{genesis.Hash(): genesis.Difficulty()},
+ ownCanonical: map[uint64]common.Hash{0: genesis.Hash()},
}
tester.stateDb = rawdb.NewMemoryDatabase()
tester.triedb = trie.NewDatabase(tester.stateDb)
@@ -157,9 +173,13 @@ func (dl *downloadTester) HasHeader(hash common.Hash, number uint64) bool {
return dl.GetHeaderByHash(hash) != nil
}
-// HasBlock checks if a block is present in the testers canonical chain.
+// HasBlock checks whether a full block is present at the requested height.
func (dl *downloadTester) HasBlock(hash common.Hash, number uint64) bool {
- return dl.GetBlockByHash(hash) != nil
+ dl.lock.RLock()
+ defer dl.lock.RUnlock()
+
+ block := dl.ownBlocks[hash]
+ return block != nil && block.NumberU64() == number
}
// HasFastBlock checks if a block is present in the testers canonical chain.
@@ -282,7 +302,12 @@ func (dl *downloadTester) InsertHeaderChain(headers []*types.Header, checkFreq i
hashes = append(hashes, hash)
}
hashes = append(hashes, headers[len(headers)-1].Hash())
+ // Anchor of the WriteHeader comparison: the current head header. A nil
+ // total difficulty means the entry was rolled back.
+ _, headHash := dl.canonicalHead(false)
+ headTd := dl.ownChainTd[headHash]
// Do a full insert if pre-checks passed
+ tookOver := false
for i, header := range headers {
hash := hashes[i]
if dl.getHeaderByHash(hash) != nil {
@@ -292,11 +317,26 @@ func (dl *downloadTester) InsertHeaderChain(headers []*types.Header, checkFreq i
// This _should_ be impossible, due to precheck and induction
return i, fmt.Errorf("InsertHeaderChain: unknown parent at position %d", i)
}
- dl.ownHashes = append(dl.ownHashes, hash)
dl.ownHeaders[hash] = header
td := dl.getTd(header.ParentHash)
dl.ownChainTd[hash] = new(big.Int).Add(td, header.Difficulty)
+
+ // Mirror WriteHeader: a header with a higher total difficulty than
+ // the head takes over the canonical table, before any body exists.
+ // WriteHeader's equal-total-difficulty clause splits ties randomly
+ // (mrand.Float64() < 0.5), which cannot be mirrored deterministically;
+ // this tester keeps the strict comparison, the conservative side of
+ // that distribution.
+ if headTd == nil || dl.ownChainTd[hash].Cmp(headTd) > 0 {
+ dl.canonicalize(hash, header.Number.Uint64())
+ dl.removeAbove(header.Number.Uint64())
+ headTd = dl.ownChainTd[hash]
+ tookOver = true
+ }
+ }
+ if tookOver {
+ dl.rebuildOwnHashes()
}
return len(headers), nil
}
@@ -306,19 +346,29 @@ func (dl *downloadTester) InsertChain(blocks types.Blocks) (i int, err error) {
dl.lock.Lock()
defer dl.lock.Unlock()
+ if dl.parkTailOnce {
+ dl.parkTailOnce = false
+ // Park the tail the way insertChain does with future blocks.
+ for i, block := range blocks[:len(blocks)-1] {
+ if err := dl.storeBlock(block); err != nil {
+ return i, fmt.Errorf("InsertChain: %v at position %d / %d", err, i, len(blocks))
+ }
+ }
+ return len(blocks) - 1, nil
+ }
+
for i, block := range blocks {
- if parent, ok := dl.ownBlocks[block.ParentHash()]; !ok {
- return i, fmt.Errorf("InsertChain: unknown parent at position %d / %d", i, len(blocks))
- } else if _, err := dl.stateDb.Get(parent.Root().Bytes()); err != nil {
- return i, fmt.Errorf("InsertChain: unknown parent state %x: %v", parent.Root(), err)
+ if err := dl.storeBlock(block); err != nil {
+ return i, fmt.Errorf("InsertChain: %v at position %d / %d", err, i, len(blocks))
}
- if _, ok := dl.ownHeaders[block.Hash()]; !ok {
- dl.ownHashes = append(dl.ownHashes, block.Hash())
- dl.ownHeaders[block.Hash()] = block.Header()
+ }
+ if hook := dl.extendTailAfterInsert; hook != nil {
+ dl.extendTailAfterInsert = nil
+ if child := hook(blocks[len(blocks)-1]); child != nil {
+ if err := dl.storeBlock(child); err != nil {
+ return len(blocks), fmt.Errorf("InsertChain: %v at position %d / %d", err, len(blocks), len(blocks)+1)
+ }
}
- dl.ownBlocks[block.Hash()] = block
- dl.stateDb.Put(block.Root().Bytes(), []byte{0x00})
- dl.ownChainTd[block.Hash()] = new(big.Int).Add(dl.ownChainTd[block.ParentHash()], block.Difficulty())
}
return len(blocks), nil
}
@@ -347,6 +397,135 @@ func (dl *downloadTester) writeBlockWithoutState(block *types.Block) error {
return nil
}
+// storeBlock records a block in the tester's chain, validating its parent.
+// It must be called with dl.lock held for writing.
+func (dl *downloadTester) storeBlock(block *types.Block) error {
+ parent, ok := dl.ownBlocks[block.ParentHash()]
+ if !ok {
+ return fmt.Errorf("unknown parent %s", block.ParentHash())
+ }
+ if _, err := dl.stateDb.Get(parent.Root().Bytes()); err != nil {
+ return fmt.Errorf("unknown parent state %x: %v", parent.Root(), err)
+ }
+ dl.ownHeaders[block.Hash()] = block.Header()
+ dl.ownBlocks[block.Hash()] = block
+ dl.stateDb.Put(block.Root().Bytes(), []byte{0x00})
+ dl.ownChainTd[block.Hash()] = new(big.Int).Add(dl.ownChainTd[block.ParentHash()], block.Difficulty())
+ // Mirror insertChain: only a block heavier than the head block takes
+ // over the canonical table; a lighter one is written as a side entry.
+ // On an equal total difficulty, production splits the tie by number:
+ // a higher-height block takes over, a same-height one stays a side
+ // entry (the selfish-mining guard of writeBlockWithState).
+ headNumber, headHash := dl.canonicalHead(true)
+ headTd := dl.ownChainTd[headHash]
+ blockTd := dl.ownChainTd[block.Hash()]
+ if headTd == nil || blockTd.Cmp(headTd) > 0 || (blockTd.Cmp(headTd) == 0 && block.NumberU64() > headNumber) {
+ dl.canonicalize(block.Hash(), block.NumberU64())
+ // Mirror the reorg cleanup of real insertChain: when a shorter but
+ // heavier branch takes over, the stale markers above the new head
+ // must go, so GetCanonicalHash stops answering with the replaced
+ // branch at those heights.
+ dl.removeAbove(block.NumberU64())
+ // ownHashes is the canonical snapshot, so refresh it only after
+ // removeAbove: rebuilding earlier would re-enter the replaced
+ // branch's stale markers still present in the table.
+ dl.rebuildOwnHashes()
+ }
+ return nil
+}
+
+// canonicalize makes hash the canonical entry at number and rewires the
+// segment below it, mirroring the marker writes of WriteHeader and
+// insertChain. Caller must hold dl.lock for writing.
+func (dl *downloadTester) canonicalize(hash common.Hash, number uint64) {
+ for n, h := number, hash; n > 0 && dl.ownCanonical[n] != h; n-- {
+ dl.ownCanonical[n] = h
+ header := dl.ownHeaders[h]
+ if header == nil {
+ break
+ }
+ h = header.ParentHash
+ }
+}
+
+// removeAbove drops canonical entries above number, the stale-marker cleanup
+// of WriteHeader. Caller must hold dl.lock for writing.
+func (dl *downloadTester) removeAbove(number uint64) {
+ for n := number + 1; ; n++ {
+ if _, ok := dl.ownCanonical[n]; !ok {
+ break
+ }
+ delete(dl.ownCanonical, n)
+ }
+}
+
+// rebuildOwnHashes refreshes ownHashes from the canonical table, keeping it
+// the height-ordered snapshot of the canonical chain that the head getters
+// report. It must run after removeAbove, so the stale markers of a replaced
+// branch never re-enter the list. Caller must hold dl.lock for writing.
+func (dl *downloadTester) rebuildOwnHashes() {
+ numbers := make([]uint64, 0, len(dl.ownCanonical))
+ for n, h := range dl.ownCanonical {
+ if dl.ownHeaders[h] != nil {
+ numbers = append(numbers, n)
+ }
+ }
+ slices.Sort(numbers)
+ dl.ownHashes = dl.ownHashes[:0]
+ for _, n := range numbers {
+ dl.ownHashes = append(dl.ownHashes, dl.ownCanonical[n])
+ }
+}
+
+// canonicalHead returns the highest canonical entry. With requireBlock,
+// entries without a stored block are skipped, giving the head block that
+// insertChain reorganizes against; without it, the head header that
+// WriteHeader compares against. Caller must hold dl.lock.
+func (dl *downloadTester) canonicalHead(requireBlock bool) (uint64, common.Hash) {
+ headNum, headHash := uint64(0), dl.ownCanonical[0]
+ for n, h := range dl.ownCanonical {
+ if n > headNum && (!requireBlock || dl.ownBlocks[h] != nil) {
+ headNum, headHash = n, h
+ }
+ }
+ return headNum, headHash
+}
+
+// GetCanonicalHash returns the hash of the canonical block at the given
+// height, or the zero hash when the height is unknown. Entries are advanced
+// by header and block inserts that overtake the head's total difficulty,
+// mirroring WriteHeader and insertChain.
+func (dl *downloadTester) GetCanonicalHash(number uint64) common.Hash {
+ dl.lock.RLock()
+ defer dl.lock.RUnlock()
+
+ return dl.ownCanonical[number]
+}
+
+// GetHeaderByNumber returns the canonical header at the given height, or nil
+// when the height is unknown, mirroring the real chain's canonical marker
+// reads used by consensus.ShouldHandleProposedBlock.
+func (dl *downloadTester) GetHeaderByNumber(number uint64) *types.Header {
+ dl.lock.RLock()
+ defer dl.lock.RUnlock()
+
+ return dl.ownHeaders[dl.ownCanonical[number]]
+}
+
+// GetBlock returns the stored block with the requested hash and height, or
+// nil when it is not stored, mirroring the real chain's storage check used by
+// consensus.ShouldHandleProposedBlock.
+func (dl *downloadTester) GetBlock(hash common.Hash, number uint64) *types.Block {
+ dl.lock.RLock()
+ defer dl.lock.RUnlock()
+
+ block := dl.ownBlocks[hash]
+ if block == nil || block.NumberU64() != number {
+ return nil
+ }
+ return block
+}
+
// InsertReceiptChain injects a new batch of receipts into the simulated chain.
func (dl *downloadTester) InsertReceiptChain(blocks types.Blocks, receipts []types.Receipts) (i int, err error) {
dl.lock.Lock()
@@ -371,14 +550,22 @@ func (dl *downloadTester) Rollback(hashes []common.Hash) {
defer dl.lock.Unlock()
for i := len(hashes) - 1; i >= 0; i-- {
- if dl.ownHashes[len(dl.ownHashes)-1] == hashes[i] {
- dl.ownHashes = dl.ownHashes[:len(dl.ownHashes)-1]
+ // Drop the canonical marker as well, otherwise GetCanonicalHash
+ // keeps reporting blocks that Rollback just removed. The height is
+ // read before the header itself is deleted below.
+ if header := dl.ownHeaders[hashes[i]]; header != nil {
+ if number := header.Number.Uint64(); dl.ownCanonical[number] == hashes[i] {
+ delete(dl.ownCanonical, number)
+ }
}
delete(dl.ownChainTd, hashes[i])
delete(dl.ownHeaders, hashes[i])
delete(dl.ownReceipts, hashes[i])
delete(dl.ownBlocks, hashes[i])
}
+ // ownHashes is the canonical snapshot, so rebuild it from the markers
+ // the loop above may have dropped.
+ dl.rebuildOwnHashes()
}
// newPeer registers a new block download source into the downloader.
@@ -400,11 +587,23 @@ func (dl *downloadTester) dropPeer(id string) {
dl.downloader.UnregisterPeer(id)
}
-// an empty handleProposedBlock function
+// handleProposedBlock records the invocation and its argument for assertions.
func (dl *downloadTester) handleProposedBlock(header *types.Header) error {
+ dl.lock.Lock()
+ dl.proposedCalls++
+ dl.lastProposedHeader = header
+ dl.lock.Unlock()
return nil
}
+// proposedState returns the invocation count and the header of the most
+// recent handleProposedBlock call.
+func (dl *downloadTester) proposedState() (int, *types.Header) {
+ dl.lock.RLock()
+ defer dl.lock.RUnlock()
+ return dl.proposedCalls, dl.lastProposedHeader
+}
+
// Config retrieves the blockchain's chain configuration.
func (dl *downloadTester) Config() *params.ChainConfig {
if dl.configOverride != nil {
@@ -2588,6 +2787,10 @@ func testReorgProtectionDoesNotStallSync(t *testing.T, protocol int, mode SyncMo
tester.ownHashes = append(tester.ownHashes[:0], localChain.chain...)
for hash, header := range localChain.headerm {
tester.ownHeaders[hash] = header
+ // Mirror the canonical markers as well: rebuildOwnHashes sources
+ // from ownCanonical, and the first canonicalize of the sync
+ // would otherwise wipe the preset chain back to genesis.
+ tester.ownCanonical[header.Number.Uint64()] = hash
}
for _, block := range localChain.blockm {
tester.ownBlocks[block.Hash()] = block
@@ -3018,3 +3221,493 @@ func TestDownloaderUnregisterPeerNeverRegistered(t *testing.T) {
t.Fatalf("unregistering a never-registered peer error mismatch: got %v want %v", err, errNotRegistered)
}
}
+
+// makeTestBlockWithDifficulty builds an empty block on top of parent with a
+// custom difficulty, so a shorter branch can out-cumulative-difficulty a
+// longer one the way production reorgs are decided. The seed goes into the
+// extra data, so chains from the same parent with different seeds are
+// distinct forks.
+func makeTestBlockWithDifficulty(parent *types.Block, difficulty int64, seed byte) *types.Block {
+ header := &types.Header{
+ ParentHash: parent.Hash(),
+ Number: new(big.Int).Add(parent.Number(), common.Big1),
+ Difficulty: big.NewInt(difficulty),
+ GasLimit: params.GenesisGasLimit,
+ Time: parent.Header().Time + 10,
+ Extra: []byte{seed},
+ }
+ return types.NewBlockWithHeader(header).WithBody(types.Body{})
+}
+
+// makeTestBlock builds an empty block on top of parent with the default
+// difficulty of one.
+func makeTestBlock(parent *types.Block, seed byte) *types.Block {
+ return makeTestBlockWithDifficulty(parent, common.Big1.Int64(), seed)
+}
+
+// extendTestChain builds n blocks on top of parent, seeding their extra data
+// with seed, seed+1 and so on. parent is not included.
+func extendTestChain(parent *types.Block, n int, seed byte) []*types.Block {
+ chain := make([]*types.Block, 0, n)
+ for i := 0; i < n; i++ {
+ block := makeTestBlock(parent, seed+byte(i))
+ chain = append(chain, block)
+ parent = block
+ }
+ return chain
+}
+
+// assertProposedBlock checks the handleProposedBlock call count and, when
+// wantTail is non-nil, the block the most recent call saw.
+func assertProposedBlock(t *testing.T, dl *downloadTester, wantCalls int, wantTail *types.Block) {
+ t.Helper()
+ calls, got := dl.proposedState()
+ if calls != wantCalls {
+ t.Fatalf("handleProposedBlock ran %d times, want %d", calls, wantCalls)
+ }
+ if wantTail == nil {
+ return
+ }
+ if got == nil || got.Hash() != wantTail.Hash() {
+ t.Fatalf("handleProposedBlock ran on %v, want %v", got, wantTail.Hash())
+ }
+}
+
+// TestImportBlockResultsProposedBlockHandler checks that the handler only
+// runs for a batch tail that is both stored and canonical, since a nil
+// InsertChain also covers a parked tail, a fork batch stored as side
+// entries, and a fast-sync height made canonical before its body arrives.
+func TestImportBlockResultsProposedBlockHandler(t *testing.T) {
+ t.Run("batch tail parked in the future queue", func(t *testing.T) {
+ // A parked tail is neither stored nor canonical, so the handler must
+ // not run.
+ dl := newTester()
+ defer dl.terminate()
+
+ local := extendTestChain(dl.genesis, 4, 0)
+ if _, err := dl.InsertChain(local); err != nil {
+ t.Fatalf("failed to set up the local chain: %v", err)
+ }
+ // The batch continues the local head, like a queued batch does.
+ batch := extendTestChain(local[len(local)-1], 4, 16)
+ tail := batch[len(batch)-1]
+ dl.parkTailOnce = true
+ before, _ := dl.proposedState()
+ if err := dl.downloader.importBlockResults(toFetchResults(batch)); err != nil {
+ t.Fatalf("failed to import the queued batch: %v", err)
+ }
+ assertProposedBlock(t, dl, before, nil)
+ if dl.GetBlock(tail.Hash(), tail.NumberU64()) != nil {
+ t.Fatalf("queued batch tail unexpectedly stored")
+ }
+ if got := dl.GetCanonicalHash(tail.NumberU64()); got != (common.Hash{}) {
+ t.Fatalf("parked tail height unexpectedly canonical: %v", got)
+ }
+ })
+
+ t.Run("imported batch", func(t *testing.T) {
+ // A fully written batch triggers the handler exactly once.
+ dl := newTester()
+ defer dl.terminate()
+
+ batch := extendTestChain(dl.genesis, 4, 32)
+ before, _ := dl.proposedState()
+ if err := dl.downloader.importBlockResults(toFetchResults(batch)); err != nil {
+ t.Fatalf("failed to import the batch: %v", err)
+ }
+ assertProposedBlock(t, dl, before+1, batch[len(batch)-1])
+ if dl.GetBlock(batch[len(batch)-1].Hash(), batch[len(batch)-1].NumberU64()) == nil {
+ t.Fatalf("imported batch tail not stored")
+ }
+ })
+
+ t.Run("stored side-chain batch re-delivered", func(t *testing.T) {
+ // A stored fork batch that is no longer the head must not reach the
+ // handler: storage is not canonicality.
+ dl := newTester()
+ defer dl.terminate()
+
+ fork := extendTestChain(dl.genesis, 4, 48)
+ // Freshly stored, the fork is ahead of the local chain, so it
+ // becomes the head and the handler firing once is expected.
+ if err := dl.downloader.importBlockResults(toFetchResults(fork)); err != nil {
+ t.Fatalf("failed to store the fork batch: %v", err)
+ }
+ // Growing the local chain past the fork turns its blocks into side
+ // entries.
+ local := extendTestChain(dl.genesis, 5, 64)
+ if _, err := dl.InsertChain(local); err != nil {
+ t.Fatalf("failed to set up the local chain: %v", err)
+ }
+ tail := fork[len(fork)-1]
+ if dl.GetBlock(tail.Hash(), tail.NumberU64()) == nil {
+ t.Fatalf("fork batch tail not stored")
+ }
+ // Re-delivering the stored fork must not reach the handler again.
+ before, _ := dl.proposedState()
+ if err := dl.downloader.importBlockResults(toFetchResults(fork)); err != nil {
+ t.Fatalf("failed to re-import the fork batch: %v", err)
+ }
+ assertProposedBlock(t, dl, before, nil)
+ })
+
+ t.Run("head advanced past the canonical tail", func(t *testing.T) {
+ // The tail stays canonical at its own height even once the head
+ // moved past it, so the handler must still run.
+ dl := newTester()
+ defer dl.terminate()
+
+ base := extendTestChain(dl.genesis, 4, 80)
+ if _, err := dl.InsertChain(base); err != nil {
+ t.Fatalf("failed to set up the local chain: %v", err)
+ }
+ batch := extendTestChain(base[len(base)-1], 4, 96)
+ dl.extendTailAfterInsert = func(tail *types.Block) *types.Block {
+ // A concurrent import of the tail's child landing before
+ // importBlockResults reads the head.
+ return makeTestBlock(tail, 200)
+ }
+ before, _ := dl.proposedState()
+ if err := dl.downloader.importBlockResults(toFetchResults(batch)); err != nil {
+ t.Fatalf("failed to import the batch: %v", err)
+ }
+ assertProposedBlock(t, dl, before+1, batch[len(batch)-1])
+ })
+
+ t.Run("heavier fork stays canonical", func(t *testing.T) {
+ // The head goes to the highest total difficulty, not to the last
+ // insert, so a lighter chain stored later must not displace the
+ // heavier fork.
+ dl := newTester()
+ defer dl.terminate()
+
+ fork := extendTestChain(dl.genesis, 6, 112)
+ if err := dl.downloader.importBlockResults(toFetchResults(fork)); err != nil {
+ t.Fatalf("failed to store the fork batch: %v", err)
+ }
+ local := extendTestChain(dl.genesis, 2, 128)
+ if _, err := dl.InsertChain(local); err != nil {
+ t.Fatalf("failed to set up the local chain: %v", err)
+ }
+ // The fork is still canonical, so re-delivering it reaches the
+ // handler again.
+ before, _ := dl.proposedState()
+ if err := dl.downloader.importBlockResults(toFetchResults(fork)); err != nil {
+ t.Fatalf("failed to re-import the fork batch: %v", err)
+ }
+ assertProposedBlock(t, dl, before+1, fork[len(fork)-1])
+ })
+
+ t.Run("fast-sync canonical tail without body", func(t *testing.T) {
+ // The header phase makes the tail's height canonical before its body
+ // is imported, so HasBlock has to gate the handler too.
+ dl := newTester()
+ defer dl.terminate()
+
+ batch := extendTestChain(dl.genesis, 4, 144)
+ headers := make([]*types.Header, len(batch))
+ for i, block := range batch {
+ headers[i] = block.Header()
+ }
+ if _, err := dl.InsertHeaderChain(headers, 0); err != nil {
+ t.Fatalf("failed to set up the header phase: %v", err)
+ }
+ // The tail height is canonical already, but its body is missing.
+ tail := batch[len(batch)-1]
+ if got := dl.GetCanonicalHash(tail.NumberU64()); got != tail.Hash() {
+ t.Fatalf("tail height not canonical after the header phase: have %v, want %v", got, tail.Hash())
+ }
+ if dl.HasBlock(tail.Hash(), tail.NumberU64()) {
+ t.Fatalf("tail block unexpectedly stored before the body phase")
+ }
+ dl.parkTailOnce = true
+ before, _ := dl.proposedState()
+ if err := dl.downloader.importBlockResults(toFetchResults(batch)); err != nil {
+ t.Fatalf("failed to import the queued batch: %v", err)
+ }
+ assertProposedBlock(t, dl, before, nil)
+ if dl.HasBlock(tail.Hash(), tail.NumberU64()) {
+ t.Fatalf("queued batch tail unexpectedly stored")
+ }
+ })
+}
+
+// TestInsertChainErrorReportsPosition checks that a failing storeBlock is
+// reported with the block's index in the batch and the batch length, so a
+// broken batch can be located without re-deriving the position.
+func TestInsertChainErrorReportsPosition(t *testing.T) {
+ dl := newTester()
+ defer dl.terminate()
+
+ base := extendTestChain(dl.genesis, 2, 208)
+ // A child hanging off a parent that is not in the chain.
+ orphanChild := makeTestBlock(makeTestBlock(dl.genesis, 250), 251)
+
+ t.Run("main loop", func(t *testing.T) {
+ batch := []*types.Block{base[0], orphanChild}
+ i, err := dl.InsertChain(batch)
+ if err == nil {
+ t.Fatalf("expected an error for the unknown parent")
+ }
+ want := fmt.Sprintf("InsertChain: unknown parent %s at position 1 / 2", orphanChild.ParentHash())
+ if err.Error() != want {
+ t.Fatalf("error mismatch: have %q, want %q", err.Error(), want)
+ }
+ if i != 1 {
+ t.Fatalf("returned index mismatch: have %d, want 1", i)
+ }
+ })
+
+ t.Run("parked tail prefix", func(t *testing.T) {
+ dl.parkTailOnce = true
+ batch := []*types.Block{orphanChild, base[1]}
+ i, err := dl.InsertChain(batch)
+ if err == nil {
+ t.Fatalf("expected an error for the unknown parent")
+ }
+ want := fmt.Sprintf("InsertChain: unknown parent %s at position 0 / 2", orphanChild.ParentHash())
+ if err.Error() != want {
+ t.Fatalf("error mismatch: have %q, want %q", err.Error(), want)
+ }
+ if i != 0 {
+ t.Fatalf("returned index mismatch: have %d, want 0", i)
+ }
+ })
+
+ t.Run("extended tail child", func(t *testing.T) {
+ // The hook's child fails on its unknown parent, reported at the
+ // position right after the batch.
+ dl.extendTailAfterInsert = func(tail *types.Block) *types.Block {
+ return orphanChild
+ }
+ batch := []*types.Block{base[0], base[1]}
+ i, err := dl.InsertChain(batch)
+ if err == nil {
+ t.Fatalf("expected an error for the unknown parent")
+ }
+ want := fmt.Sprintf("InsertChain: unknown parent %s at position 2 / 3", orphanChild.ParentHash())
+ if err.Error() != want {
+ t.Fatalf("error mismatch: have %q, want %q", err.Error(), want)
+ }
+ if i != 2 {
+ t.Fatalf("returned index mismatch: have %d, want 2", i)
+ }
+ })
+}
+
+// TestRollbackClearsCanonicalMarkers checks that Rollback drops the
+// canonical entries of the removed blocks, so GetCanonicalHash stops
+// reporting them, while the surviving ancestors stay canonical.
+func TestRollbackClearsCanonicalMarkers(t *testing.T) {
+ dl := newTester()
+ defer dl.terminate()
+
+ chain := extendTestChain(dl.genesis, 4, 208)
+ if _, err := dl.InsertChain(chain); err != nil {
+ t.Fatalf("failed to set up the chain: %v", err)
+ }
+ for _, block := range chain {
+ if got := dl.GetCanonicalHash(block.NumberU64()); got != block.Hash() {
+ t.Fatalf("height %d not canonical before the rollback: have %v, want %v", block.NumberU64(), got, block.Hash())
+ }
+ }
+
+ rolledBack := chain[2:]
+ hashes := make([]common.Hash, len(rolledBack))
+ for i, block := range rolledBack {
+ hashes[i] = block.Hash()
+ }
+ dl.Rollback(hashes)
+
+ // The rolled-back heights no longer report a canonical hash.
+ for _, block := range rolledBack {
+ if got := dl.GetCanonicalHash(block.NumberU64()); got != (common.Hash{}) {
+ t.Fatalf("height %d still canonical after the rollback: %v", block.NumberU64(), got)
+ }
+ }
+ // The surviving ancestors keep their canonical entries.
+ for _, block := range chain[:2] {
+ if got := dl.GetCanonicalHash(block.NumberU64()); got != block.Hash() {
+ t.Fatalf("ancestor height %d lost its canonical entry: have %v, want %v", block.NumberU64(), got, block.Hash())
+ }
+ }
+ if got := dl.GetCanonicalHash(0); got != dl.genesis.Hash() {
+ t.Fatalf("genesis height lost its canonical entry: have %v, want %v", got, dl.genesis.Hash())
+ }
+}
+
+// toFetchResults converts a chain of blocks into the fetch-result shape
+// importBlockResults consumes: headers plus empty transaction bodies.
+func toFetchResults(blocks []*types.Block) []*fetchResult {
+ results := make([]*fetchResult, 0, len(blocks))
+ for _, block := range blocks {
+ results = append(results, &fetchResult{Header: block.Header(), Transactions: types.Transactions{}})
+ }
+ return results
+}
+
+// TestStoreBlockCleansStaleCanonicalMarkers checks that taking over the
+// canonical table with a shorter but heavier branch drops the stale
+// canonical markers above the new head, mirroring the reorg cleanup of
+// real insertChain: without the cleanup, GetCanonicalHash keeps answering
+// with the replaced branch at the heights above the new head.
+func TestStoreBlockCleansStaleCanonicalMarkers(t *testing.T) {
+ dl := newTester()
+ defer dl.terminate()
+
+ base := extendTestChain(dl.genesis, 2, 208)
+ if _, err := dl.InsertChain(base); err != nil {
+ t.Fatalf("failed to set up the base chain: %v", err)
+ }
+ baseHead := base[len(base)-1]
+
+ // A longer branch of difficulty-1 blocks: heights baseHead+1..+3.
+ longBranch := extendTestChain(baseHead, 3, 226)
+ if _, err := dl.InsertChain(longBranch); err != nil {
+ t.Fatalf("failed to import the long branch: %v", err)
+ }
+
+ // A shorter branch of difficulty-2 blocks: heights baseHead+1..+2
+ // only, but its head carries the higher total difficulty.
+ var shortBranch []*types.Block
+ parent := baseHead
+ for i := 0; i < 2; i++ {
+ parent = makeTestBlockWithDifficulty(parent, 2, 240+byte(i))
+ shortBranch = append(shortBranch, parent)
+ }
+ if _, err := dl.InsertChain(shortBranch); err != nil {
+ t.Fatalf("failed to import the short heavier branch: %v", err)
+ }
+
+ // The short branch took the canonical table over up to its head.
+ head := shortBranch[len(shortBranch)-1]
+ if got := dl.GetCanonicalHash(head.NumberU64()); got != head.Hash() {
+ t.Fatalf("heavier fork head not canonical: have %v, want %v", got, head.Hash())
+ }
+ if got := dl.GetCanonicalHash(shortBranch[0].NumberU64()); got != shortBranch[0].Hash() {
+ t.Fatalf("replaced height %d lost its new canonical entry: have %v, want %v", shortBranch[0].NumberU64(), got, shortBranch[0].Hash())
+ }
+ // The heights above the new head must not answer with the replaced
+ // branch; production reorg deletes those markers.
+ above := longBranch[len(longBranch)-1]
+ if got := dl.GetCanonicalHash(above.NumberU64()); got != (common.Hash{}) {
+ t.Fatalf("stale canonical marker above the new head survived: have %v", got)
+ }
+}
+
+// TestStoreBlockEqualDifficultyTieBreaksByNumber checks that the fork
+// choice mirror of insertChain splits an equal-total-difficulty tie by
+// number: a higher-height block takes over the canonical table while a
+// same-height one stays a side entry, matching the selfish-mining guard
+// of writeBlockWithState.
+func TestStoreBlockEqualDifficultyTieBreaksByNumber(t *testing.T) {
+ dl := newTester()
+ defer dl.terminate()
+
+ base := extendTestChain(dl.genesis, 2, 250)
+ if _, err := dl.InsertChain(base); err != nil {
+ t.Fatalf("failed to set up the base chain: %v", err)
+ }
+ baseHead := base[len(base)-1]
+
+ // A difficulty-2 block at height baseHead+1 takes the head with a
+ // strictly higher total difficulty.
+ heavy := makeTestBlockWithDifficulty(baseHead, 2, 251)
+ if _, err := dl.InsertChain([]*types.Block{heavy}); err != nil {
+ t.Fatalf("failed to import the heavy fork: %v", err)
+ }
+ if got := dl.GetCanonicalHash(heavy.NumberU64()); got != heavy.Hash() {
+ t.Fatalf("heavier fork head not canonical: have %v, want %v", got, heavy.Hash())
+ }
+
+ // A same-height block with the same total difficulty must not
+ // displace the head.
+ twin := makeTestBlockWithDifficulty(baseHead, 2, 252)
+ if _, err := dl.InsertChain([]*types.Block{twin}); err != nil {
+ t.Fatalf("failed to import the twin fork: %v", err)
+ }
+ if got := dl.GetCanonicalHash(twin.NumberU64()); got != heavy.Hash() {
+ t.Fatalf("same-height equal-difficulty fork displaced the head: have %v, want %v", got, heavy.Hash())
+ }
+
+ // A branch of difficulty-1 blocks ties the head's total difficulty
+ // exactly one height above it; production promotes that block, so
+ // the mirror must too.
+ tie := extendTestChain(baseHead, 2, 253)
+ if _, err := dl.InsertChain(tie); err != nil {
+ t.Fatalf("failed to import the tying branch: %v", err)
+ }
+ if got := dl.GetCanonicalHash(tie[1].NumberU64()); got != tie[1].Hash() {
+ t.Fatalf("equal-difficulty higher-height fork not promoted: have %v, want %v", got, tie[1].Hash())
+ }
+ if got := dl.GetCanonicalHash(tie[0].NumberU64()); got != tie[0].Hash() {
+ t.Fatalf("tying branch lost its canonical segment: have %v, want %v", got, tie[0].Hash())
+ }
+ if got := dl.GetCanonicalHash(baseHead.NumberU64()); got != baseHead.Hash() {
+ t.Fatalf("reorg rewired the shared prefix: have %v, want %v", got, baseHead.Hash())
+ }
+}
+
+// TestHeadGettersFollowCanonicalTable checks that the head getters report
+// the canonical head rather than the last stored block, the way production
+// reads its own chain state: a lighter side chain stored after a heavier
+// fork must not displace what CurrentHeader, CurrentBlock and
+// CurrentSnapBlock report, and a shorter heavier branch taking over must
+// move the getters onto its own head, off the replaced branch's tail.
+func TestHeadGettersFollowCanonicalTable(t *testing.T) {
+ dl := newTester()
+ defer dl.terminate()
+
+ fork := extendTestChain(dl.genesis, 6, 112)
+ if err := dl.downloader.importBlockResults(toFetchResults(fork)); err != nil {
+ t.Fatalf("failed to store the fork batch: %v", err)
+ }
+ local := extendTestChain(dl.genesis, 2, 128)
+ if _, err := dl.InsertChain(local); err != nil {
+ t.Fatalf("failed to set up the local chain: %v", err)
+ }
+ want := fork[len(fork)-1].NumberU64()
+ if have := dl.CurrentHeader().Number.Uint64(); have != want {
+ t.Fatalf("CurrentHeader reported the side chain: have %v, want %v", have, want)
+ }
+ if have := dl.CurrentBlock().Number.Uint64(); have != want {
+ t.Fatalf("CurrentBlock reported the side chain: have %v, want %v", have, want)
+ }
+ if have := dl.CurrentSnapBlock().Number.Uint64(); have != want {
+ t.Fatalf("CurrentSnapBlock reported the side chain: have %v, want %v", have, want)
+ }
+
+ // A shorter heavier branch taking over the canonical table must move the
+ // head getters onto its own head and off the replaced branch's stale tail.
+ dl2 := newTester()
+ defer dl2.terminate()
+
+ base := extendTestChain(dl2.genesis, 2, 208)
+ if _, err := dl2.InsertChain(base); err != nil {
+ t.Fatalf("failed to set up the base chain: %v", err)
+ }
+ baseHead := base[len(base)-1]
+ longBranch := extendTestChain(baseHead, 3, 226)
+ if _, err := dl2.InsertChain(longBranch); err != nil {
+ t.Fatalf("failed to import the long branch: %v", err)
+ }
+ var shortBranch []*types.Block
+ parent := baseHead
+ for i := 0; i < 2; i++ {
+ parent = makeTestBlockWithDifficulty(parent, 2, 240+byte(i))
+ shortBranch = append(shortBranch, parent)
+ }
+ if _, err := dl2.InsertChain(shortBranch); err != nil {
+ t.Fatalf("failed to import the short heavier branch: %v", err)
+ }
+ want = shortBranch[len(shortBranch)-1].NumberU64()
+ if have := dl2.CurrentHeader().Number.Uint64(); have != want {
+ t.Fatalf("CurrentHeader kept the replaced branch: have %v, want %v", have, want)
+ }
+ if have := dl2.CurrentBlock().Number.Uint64(); have != want {
+ t.Fatalf("CurrentBlock kept the replaced branch: have %v, want %v", have, want)
+ }
+ if have := dl2.CurrentSnapBlock().Number.Uint64(); have != want {
+ t.Fatalf("CurrentSnapBlock kept the replaced branch: have %v, want %v", have, want)
+ }
+}
diff --git a/eth/fetcher/block_fetcher.go b/eth/fetcher/block_fetcher.go
index caea7b6a3483..6216c020f310 100644
--- a/eth/fetcher/block_fetcher.go
+++ b/eth/fetcher/block_fetcher.go
@@ -223,6 +223,15 @@ func NewBlockFetcher(getBlock blockRetrievalFn, verifyHeader headerVerifierFn, h
}
}
+// HandleProposedBlock returns the proposed-block handler the fetcher was
+// wired with. Read-only, for tests: the eth package asserts that its gated
+// closure — not the bare consensus handler — actually reached the fetcher,
+// since the handler tests exercise the closure directly and would stay green
+// if the wiring were silently swapped.
+func (f *BlockFetcher) HandleProposedBlock() proposeBlockHandlerFn {
+ return f.handleProposedBlock
+}
+
// Start boots up the announcement based synchroniser, accepting and processing
// hash notifications and block fetches until termination requested.
func (f *BlockFetcher) Start() {
diff --git a/eth/handler.go b/eth/handler.go
index 316c7d5f0f5f..a69fa1139370 100644
--- a/eth/handler.go
+++ b/eth/handler.go
@@ -86,6 +86,11 @@ type ProtocolManager struct {
peers *peerSet
bft *bft.Bfter
+ // fetcherHandleProposedBlock is the snapSync-gated proposed-block callback
+ // the fetcher runs; kept on the manager so the fast sync gate stays
+ // observable and testable without driving the fetcher loops.
+ fetcherHandleProposedBlock func(header *types.Header) error
+
eventMux *event.TypeMux
txsCh chan core.NewTxsEvent
orderTxCh chan core.OrderTxPreEvent
@@ -175,6 +180,33 @@ func NewProtocolManager(config *params.ChainConfig, mode downloader.SyncMode, ne
}
}
+ // While fast sync runs, the fetcher must not reach the consensus handler:
+ // snapSync discards propagated blocks before executing them, so a body
+ // already written by the fast sync receipt phase would pass both halves
+ // of ShouldHandleProposedBlock and drive processQC and the vote path on
+ // a block whose state transition was never validated. Carry the same
+ // snapSync guard inserter and prepare have. The flag alone is read too
+ // late: the fetcher runs signHook between its insert and this callback,
+ // so Synchronise can flip snapSync to 0 inside that gap and let an
+ // unexecuted block through. The state check below makes that flip
+ // harmless: it is keyed by the block hash, not the state root, so it
+ // is false for any block the fetcher's inserter discarded, whatever
+ // the flag reads — a root-keyed check would leak a parent's state to
+ // an empty block whose root repeats its parent's. The downloader keeps the
+ // ungated closure on purpose: its fast sync handler calls run after the
+ // pivot commit, on blocks InsertChain has fully executed.
+ manager.fetcherHandleProposedBlock = func(header *types.Header) error {
+ if atomic.LoadUint32(&manager.snapSync) == 1 {
+ log.Debug("[fetcher] skipped proposed block handler during fast sync", "hash", header.Hash(), "number", header.Number)
+ return nil
+ }
+ if !blockchain.HasBlockAndFullState(header.Hash(), header.Number.Uint64()) {
+ log.Debug("[fetcher] skipped proposed block handler: state not executed", "hash", header.Hash(), "number", header.Number)
+ return nil
+ }
+ return handleProposedBlock(header)
+ }
+
// Construct the different synchronisation mechanisms
manager.downloader = downloader.New(chaindb, manager.eventMux, blockchain, nil, manager.removePeer, handleProposedBlock)
@@ -205,7 +237,11 @@ func NewProtocolManager(config *params.ChainConfig, mode downloader.SyncMode, ne
atomic.StoreUint32(&manager.acceptTxs, 1) // Mark initial sync done on any fetcher import
return manager.blockchain.PrepareBlock(block)
}
- manager.blockFetcher = fetcher.NewBlockFetcher(blockchain.GetBlockByHash, validator, handleProposedBlock, manager.BroadcastBlock, heighter, inserter, prepare, manager.removePeer)
+ // The fetcher must be wired with the gated closure above, not the bare
+ // handleProposedBlock: the handler tests exercise the closure directly
+ // and cannot see what the fetcher holds, so swapping this argument back
+ // would silently drop both gates.
+ manager.blockFetcher = fetcher.NewBlockFetcher(blockchain.GetBlockByHash, validator, manager.fetcherHandleProposedBlock, manager.BroadcastBlock, heighter, inserter, prepare, manager.removePeer)
fetchTx := func(peer string, hashes []common.Hash) error {
p := manager.peers.Peer(peer)
diff --git a/eth/handler_test.go b/eth/handler_test.go
index 5d40983d89ec..34952a4ebea8 100644
--- a/eth/handler_test.go
+++ b/eth/handler_test.go
@@ -21,6 +21,8 @@ import (
"math"
"math/big"
"math/rand"
+ "reflect"
+ "strings"
"sync"
"testing"
"time"
@@ -36,6 +38,7 @@ import (
"github.com/XinFinOrg/XDPoSChain/eth/downloader"
"github.com/XinFinOrg/XDPoSChain/eth/ethconfig"
"github.com/XinFinOrg/XDPoSChain/event"
+ "github.com/XinFinOrg/XDPoSChain/log"
"github.com/XinFinOrg/XDPoSChain/p2p"
"github.com/XinFinOrg/XDPoSChain/p2p/enode"
"github.com/XinFinOrg/XDPoSChain/params"
@@ -925,3 +928,84 @@ func TestRegisterDownloaderPeerUndoesRacedRemoval(t *testing.T) {
}
pm.downloader.UnregisterPeer(p.id)
}
+
+// Tests that the fetcher's proposed-block handler is inert while fast sync
+// runs. snapSync discards propagated blocks without executing them, so a
+// body already written by the fast sync receipt phase would pass both
+// halves of consensus.ShouldHandleProposedBlock and drive processQC and
+// the vote path on a block whose state transition was never validated.
+// The guard must sit in the fetcher's own callback and not in the shared
+// closure, which also feeds the downloader's post-pivot imports of fully
+// executed blocks.
+func TestFetcherSkipsProposedBlockHandlerDuringFastSync(t *testing.T) {
+ // Capture the skip log of the fetcher's callback; assert on the specific
+ // line, so unrelated background logs cannot interfere.
+ logBuf := new(lockedBuffer)
+ prevLog := log.Root()
+ glog := log.NewGlogHandler(log.NewTerminalHandlerWithLevel(logBuf, log.LevelDebug, false))
+ glog.Verbosity(log.LevelDebug)
+ log.SetDefault(log.NewLogger(glog))
+ defer log.SetDefault(prevLog)
+
+ // A fast sync manager over an empty chain keeps snapSync enabled.
+ pm, _ := newTestProtocolManagerMust(t, downloader.FastSync, 0, nil, nil)
+ defer pm.Stop()
+ if pm.snapSync != 1 {
+ t.Fatalf("fast sync not enabled: snapSync = %d", pm.snapSync)
+ }
+ // The gates are only meaningful if they actually sit in the handler the
+ // fetcher invokes. handleProposedBlock is a private fetcher field the
+ // closure-level assertions above cannot see, so read it back through the
+ // test-only accessor and require it to be the very closure defined in
+ // eth/handler.go — swapping the wiring back to the bare handleProposedBlock
+ // must fail here, not stay green.
+ wired := pm.blockFetcher.HandleProposedBlock()
+ if wired == nil {
+ t.Fatal("fetcher must be wired with a proposed-block handler")
+ }
+ if reflect.ValueOf(wired).Pointer() != reflect.ValueOf(pm.fetcherHandleProposedBlock).Pointer() {
+ t.Fatal("fetcher must be wired with the gated closure, not the bare handleProposedBlock")
+ }
+ if err := pm.fetcherHandleProposedBlock(pm.blockchain.CurrentBlock()); err != nil {
+ t.Fatalf("fetcher handler returned error during fast sync: %v", err)
+ }
+ if !strings.Contains(logBuf.String(), "skipped proposed block handler during fast sync") {
+ t.Fatalf("fast sync skip not logged, have %q", logBuf.String())
+ }
+ logBuf.Reset()
+
+ // With fast sync off, the same call reaches the closure without the skip
+ // log: the guard must not over-block full sync.
+ pmFull, _ := newTestProtocolManagerMust(t, downloader.FullSync, 0, nil, nil)
+ defer pmFull.Stop()
+ if pmFull.snapSync != 0 {
+ t.Fatalf("full sync manager must not enable snapSync: %d", pmFull.snapSync)
+ }
+ if err := pmFull.fetcherHandleProposedBlock(pmFull.blockchain.CurrentBlock()); err != nil {
+ t.Fatalf("fetcher handler returned error in full sync: %v", err)
+ }
+ if strings.Contains(logBuf.String(), "skipped proposed block handler during fast sync") {
+ t.Fatalf("full sync call must not be gated by snapSync, have %q", logBuf.String())
+ }
+ if strings.Contains(logBuf.String(), "state not executed") {
+ t.Fatalf("executed genesis block must pass the state guard, have %q", logBuf.String())
+ }
+ logBuf.Reset()
+
+ // The snapSync flag alone is read too late: the fetcher runs signHook
+ // between its insert and this callback, so Synchronise can flip the
+ // flag to 0 inside that gap. A block whose inserter was skipped never
+ // had its state executed, so the state guard must catch it even with
+ // the flag already cleared. A forged root stands in for such a block.
+ unexecuted := *pmFull.blockchain.CurrentBlock()
+ unexecuted.Root = common.HexToHash("0xdeadbeef")
+ if err := pmFull.fetcherHandleProposedBlock(&unexecuted); err != nil {
+ t.Fatalf("fetcher handler returned error for unexecuted state: %v", err)
+ }
+ if !strings.Contains(logBuf.String(), "state not executed") {
+ t.Fatalf("state guard skip not logged, have %q", logBuf.String())
+ }
+ if strings.Contains(logBuf.String(), "skipped proposed block handler during fast sync") {
+ t.Fatalf("state guard must be independent of the snapSync flag, have %q", logBuf.String())
+ }
+}