From e9e861c656f7e84baea38aa728395cb1cd4baa92 Mon Sep 17 00:00:00 2001 From: aperod-network Date: Sat, 10 Oct 2026 23:02:20 +0500 Subject: [PATCH] chore: maintain verified source distribution Co-authored-by: aperod-network --- avm/admission_policy_test.go | 74 ++++++++++++++++++++++++++++++++++ avm/block_executor.go | 7 +++- avm/validate.go | 30 +++++++++++--- core/mempool.go | 8 ++-- core/mempool_avm_nonce_test.go | 27 +++++++++++++ node/node.go | 23 +++++++++++ p2p/badblock_ban_test.go | 40 +++++++++++++++--- p2p/handshake_budget.go | 47 +++++++++++++++++++++ p2p/handshake_budget_test.go | 60 +++++++++++++++++++++++++++ p2p/host.go | 71 ++++++++++++++++++++++---------- 10 files changed, 348 insertions(+), 39 deletions(-) create mode 100644 avm/admission_policy_test.go create mode 100644 p2p/handshake_budget.go create mode 100644 p2p/handshake_budget_test.go diff --git a/avm/admission_policy_test.go b/avm/admission_policy_test.go new file mode 100644 index 0000000..05f4d65 --- /dev/null +++ b/avm/admission_policy_test.go @@ -0,0 +1,74 @@ +package avm + +import ( + "strings" + "testing" + + "github.com/aperod/aperod/core" + wabinbinary "github.com/tetratelabs/wabin/binary" + "github.com/tetratelabs/wabin/wasm" +) + +func TestInstructionAdmissionDepthBound(t *testing.T) { + for _, depth := range []int{1, MaxAdmissionNestingDepth, MaxAdmissionNestingDepth + 1} { + body := make([]byte, 0, depth*3+1) + for i := 0; i < depth; i++ { + body = append(body, 0x02, 0x40) + } + for i := 0; i <= depth; i++ { + body = append(body, 0x0b) + } + _, err := validateInstructionsBounded(body, 0, MaxAdmissionNestingDepth) + if depth <= MaxAdmissionNestingDepth && err != nil { + t.Fatal(err) + } + if depth > MaxAdmissionNestingDepth && (err == nil || !strings.Contains(err.Error(), "nesting")) { + t.Fatalf("depth %d: %v", depth, err) + } + // Historical instruction validation is deliberately unchanged. + if _, err := validateInstructions(body, 0); err != nil { + t.Fatalf("historical validator changed: %v", err) + } + } +} + +func TestModuleAdmissionUsesDepthPolicyForDeployAndCall(t *testing.T) { + for _, depth := range []int{MaxAdmissionNestingDepth, MaxAdmissionNestingDepth + 1} { + module, err := wabinbinary.DecodeModule(stateWriteModule(false), wasm.CoreFeaturesV2) + if err != nil { + t.Fatal(err) + } + prefix := make([]byte, 0, depth*3) + for i := 0; i < depth; i++ { + prefix = append(prefix, 0x02, 0x40) + } + for i := 0; i < depth; i++ { + prefix = append(prefix, 0x0b) + } + module.CodeSection[0].Body = append(prefix, module.CodeSection[0].Body...) + code := wabinbinary.EncodeModule(module) + if _, err := ValidateModule(code); err != nil { + t.Fatal(err) + } + var id [32]byte + id[0] = 1 + for _, action := range []core.AVMAction{core.AVMDeployContract, core.AVMExecuteContract} { + store := NewMemoryStore() + if action == core.AVMExecuteContract { + if err := store.Apply([]Write{{Key: contractCodeKey(id), Value: code}}); err != nil { + t.Fatal(err) + } + } + err := ValidateMempoolAdmission(store, &core.AVMPayload{ + Action: action, ContractID: id, Code: code, Entry: "run", GasLimit: 1000, + AccessList: []core.AVMAccess{{Key: []byte("key"), Write: true}}, + }) + if depth <= MaxAdmissionNestingDepth && err != nil { + t.Fatal(err) + } + if depth > MaxAdmissionNestingDepth && (err == nil || !strings.Contains(err.Error(), "nesting")) { + t.Fatalf("action %d depth %d: %v", action, depth, err) + } + } + } +} diff --git a/avm/block_executor.go b/avm/block_executor.go index 078b62d..f9f4ea3 100644 --- a/avm/block_executor.go +++ b/avm/block_executor.go @@ -29,7 +29,7 @@ type Receipt struct { // durable state. Consensus must persist Writes atomically with the block body, // canonical height index and tip. type PreparedBlock struct { - LPoD *store.LPoDSettlement + LPoD *store.LPoDSettlement Height uint64 BlockHash crypto.Hash32 Receipts []Receipt @@ -198,6 +198,11 @@ func ValidateMempoolAdmission(store Store, payload *core.AVMPayload) error { default: return fmt.Errorf("unknown action %d", payload.Action) } + // Bound compilation preflight for both submitted and previously stored code. + // Replay/consensus keeps ValidateModule's historical rules. + if _, err := validateModule(code, MaxAdmissionNestingDepth); err != nil { + return fmt.Errorf("module admission: %w", err) + } accessList := make([]Access, len(payload.AccessList)) for i, access := range payload.AccessList { accessList[i] = Access{Key: bytes.Clone(access.Key), Write: access.Write} diff --git a/avm/validate.go b/avm/validate.go index 93bda0f..d31fbf6 100644 --- a/avm/validate.go +++ b/avm/validate.go @@ -10,11 +10,12 @@ import ( ) const ( - MaxCodeSize = 1 << 20 - MaxMemoryPages = 512 // 32 MiB at 64 KiB per Wasm page. - MaxInputSize = 64 << 10 - MaxStateKeySize = 64 - MaxStateValue = 64 << 10 + MaxCodeSize = 1 << 20 + MaxMemoryPages = 512 // 32 MiB at 64 KiB per Wasm page. + MaxInputSize = 64 << 10 + MaxStateKeySize = 64 + MaxStateValue = 64 << 10 + MaxAdmissionNestingDepth = 128 // Local pool policy; not a historical consensus limit. ) var allowedImports = map[string]struct{}{ @@ -37,6 +38,10 @@ type ValidationReport struct { // indirect calls, WASI, imported memory, threads, SIMD, host pointers, time, // filesystem, network or random-number APIs. func ValidateModule(code []byte) (ValidationReport, error) { + return validateModule(code, 0) +} + +func validateModule(code []byte, maxDepth int) (ValidationReport, error) { var report ValidationReport if len(code) == 0 { return report, fmt.Errorf("avm: empty Wasm module") @@ -85,7 +90,7 @@ func ValidateModule(code []byte) (ValidationReport, error) { if hasFloatType(body.LocalTypes) { return report, fmt.Errorf("avm: function %d has floating-point locals", functionIndex) } - count, parseErr := validateInstructions(body.Body, importCount) + count, parseErr := validateInstructionsBounded(body.Body, importCount, maxDepth) if parseErr != nil { return report, fmt.Errorf("avm: function %d: %w", functionIndex, parseErr) } @@ -111,7 +116,12 @@ func isFloatType(valueType wasm.ValueType) bool { } func validateInstructions(body []byte, importCount uint32) (uint64, error) { + return validateInstructionsBounded(body, importCount, 0) +} + +func validateInstructionsBounded(body []byte, importCount uint32, maxDepth int) (uint64, error) { var count uint64 + depth := 0 for offset := 0; offset < len(body); { op := body[offset] offset++ @@ -133,11 +143,19 @@ func validateInstructions(body []byte, importCount uint32) (uint64, error) { return 0, fmt.Errorf("calls to local functions are forbidden in AVM v1") } case op == byte(wasm.OpcodeBlock), op == byte(wasm.OpcodeIf): + depth++ + if maxDepth > 0 && depth > maxDepth { + return 0, fmt.Errorf("admission nesting exceeds %d", maxDepth) + } n, err := skipBlockType(body[offset:]) if err != nil { return 0, err } offset += n + case op == byte(wasm.OpcodeEnd): + if depth > 0 { + depth-- + } case op == byte(wasm.OpcodeLocalGet), op == byte(wasm.OpcodeLocalSet), op == byte(wasm.OpcodeLocalTee), op == byte(wasm.OpcodeGlobalGet), op == byte(wasm.OpcodeGlobalSet): diff --git a/core/mempool.go b/core/mempool.go index 683d1d5..baa7cd0 100644 --- a/core/mempool.go +++ b/core/mempool.go @@ -217,10 +217,6 @@ func (m *Mempool) Add(tx Transaction) error { if err := tx.Validate(); err != nil { return fmt.Errorf("mempool: invalid tx: %w", err) } - if err := m.validateAVMNonce(&tx); err != nil { - return err - } - // Full RingCT cryptographic verification (ring sigs, range proofs, Pedersen balance). // C-0/C-1: prevents supply inflation via forged AmountCommit or unbound stake amount. // Stake txs are included — they carry ring inputs whose proofs must be valid. @@ -261,6 +257,10 @@ func (m *Mempool) Add(tx Transaction) error { } } + // Module preflight follows authorization, size and fee validation. + if err := m.validateAVMNonce(&tx); err != nil { + return err + } hash := tx.Hash() var stakeSenderKey string diff --git a/core/mempool_avm_nonce_test.go b/core/mempool_avm_nonce_test.go index 5edfda9..6f88a73 100644 --- a/core/mempool_avm_nonce_test.go +++ b/core/mempool_avm_nonce_test.go @@ -105,6 +105,33 @@ func TestMempoolAVMAdmissionCheckFailsClosed(t *testing.T) { } } +func TestMempoolAVMFeePrecedesModuleAdmission(t *testing.T) { + pool := avmNonceMempool(t, func([32]byte) (uint64, error) { return 0, nil }) + checks := 0 + pool.cfg.AVMAdmissionCheck = func(*AVMPayload) error { checks++; return nil } + tx := mempoolAVMTx(t, 0, 1) + tx.Fee = 0 + if err := pool.Add(tx); err == nil { + t.Fatal("underfunded transaction accepted") + } + if checks != 0 { + t.Fatal("module admission ran before fee validation") + } +} + +func TestMempoolAVMAuthorizationPrecedesModuleAdmission(t *testing.T) { + pool := avmNonceMempool(t, func([32]byte) (uint64, error) { return 0, nil }) + pool.cfg.Verifier = NewTxVerifier(NewUTXOSet()) + checks := 0 + pool.cfg.AVMAdmissionCheck = func(*AVMPayload) error { checks++; return nil } + if err := pool.Add(mempoolAVMTx(t, 0, 1)); err == nil { + t.Fatal("unfunded ring transaction accepted") + } + if checks != 0 { + t.Fatal("module admission ran before transaction authorization") + } +} + func TestMempoolAllowsOnlyOnePendingAVMTransactionPerSigner(t *testing.T) { pool := avmNonceMempool(t, func([32]byte) (uint64, error) { return 0, nil }) first := mempoolAVMTx(t, 0, 1) diff --git a/node/node.go b/node/node.go index f4c6cd3..91ab326 100644 --- a/node/node.go +++ b/node/node.go @@ -501,6 +501,29 @@ func (a *p2pAdapter) OnBlock(block *core.Block) { } } +func (a *p2pAdapter) ValidateBlock(block *core.Block) error { + return a.blockV.VerifyBlock(block) +} + +func (a *p2pAdapter) ValidateSyncHeader(header core.BlockHeader) error { + registry := a.engine.Registry() + if registry == nil || !registry.IsActive(header.ValidatorPub) { + return fmt.Errorf("sync header author not active") + } + authority := header.ValidatorPub + if resolver, ok := any(registry).(interface { + SignatureAuthority(crypto.ValidatorPubKey, uint64) crypto.ValidatorPubKey + }); ok { + authority = resolver.SignatureAuthority(header.ValidatorPub, header.Height) + } + signature := header.Signature + header.Signature = nil + if !authority.Verify(header.Hash(), signature) { + return fmt.Errorf("sync header signature invalid") + } + return nil +} + func (a *p2pAdapter) OnTransaction(tx *core.Transaction) { if err := a.mempool.Add(*tx); err != nil { a.log.Debug("p2p: tx rejected", "err", err) diff --git a/p2p/badblock_ban_test.go b/p2p/badblock_ban_test.go index 41d9ead..d130a02 100644 --- a/p2p/badblock_ban_test.go +++ b/p2p/badblock_ban_test.go @@ -9,6 +9,7 @@ package p2p_test // - Strike map is capped (badBlockMaxTrackedIPs) to prevent memory exhaustion. import ( + "bytes" "crypto/tls" "encoding/json" "fmt" @@ -20,6 +21,8 @@ import ( "testing" "time" + "github.com/aperod/aperod/core" + "github.com/aperod/aperod/crypto" "github.com/aperod/aperod/p2p" ) @@ -1172,11 +1175,27 @@ func (f *fixedHeightHandler) CurrentHeight() uint64 { return f.height } // 2. The test peer announces height = 980 000 in the Ping. // 3. The peer sends 20 blocks at height 980 001 (well above ourTip + 1000). // 4. After all 20 blocks the peer must NOT be banned, and PeerCount must stay 1. +type trustedSyncHandler struct { + fixedHeightHandler + author crypto.ValidatorPubKey +} + +func (h *trustedSyncHandler) ValidateSyncHeader(header core.BlockHeader) error { + if !bytes.Equal(header.ValidatorPub, h.author) || !header.VerifySignature() { + return fmt.Errorf("untrusted sync header") + } + return nil +} + func TestBadBlockBan_RelaySyncNoBan(t *testing.T) { // Relay node: local chain tip is at 1 000 (just loaded a snapshot). const relayTip = 1_000 // Validator (peer) is 979 000 blocks ahead — a realistic post-snapshot gap. const validatorHeight = 980_000 + priv, pub, err := crypto.GenerateValidatorKey() + if err != nil { + t.Fatal(err) + } h := p2p.NewHost(p2p.Config{ ListenAddr: "127.0.0.1:0", @@ -1186,7 +1205,7 @@ func TestBadBlockBan_RelaySyncNoBan(t *testing.T) { BadBlockBanThreshold: 10, BadBlockHeightLead: 1000, BadBlockBanDuration: 24 * time.Hour, - }, &fixedHeightHandler{height: relayTip}, newTestLogger()) + }, &trustedSyncHandler{fixedHeightHandler: fixedHeightHandler{height: relayTip}, author: pub}, newTestLogger()) if err := h.Start(); err != nil { t.Fatalf("Start: %v", err) } @@ -1225,11 +1244,16 @@ func TestBadBlockBan_RelaySyncNoBan(t *testing.T) { // Send 20 blocks at height 980 001 (validator's next tip block arriving // via gossip while the relay is still filling the 979 000-block gap). - // Under the old (broken) logic each block would score a rogue-fork strike. - // With the fix, none should: peer.height > ourTip + BadBlockHeightLead. + // Authenticated canonical authors remain eligible for synchronization. for i := 0; i < 20; i++ { + header := core.BlockHeader{Height: validatorHeight + 1, ValidatorPub: pub, + Timestamp: time.Now().UnixNano(), MerkleRoot: core.MerkleRoot(nil)} + if err := header.Sign(priv); err != nil { + t.Fatal(err) + } sb := p2p.SerializedBlock{ - Header: p2p.SerializedHeader{Height: validatorHeight + 1}, + Header: p2p.SerializedHeader{Height: header.Height, ValidatorPub: header.ValidatorPub[:], + Timestamp: header.Timestamp, MerkleRoot: header.MerkleRoot, Signature: header.Signature}, } conn.SetWriteDeadline(time.Now().Add(time.Second)) if err := p2p.WriteMsg(conn, p2p.MsgBlock, sb); err != nil { @@ -2298,7 +2322,9 @@ func TestSlowHandshakeFlood(t *testing.T) { // handshake completes or the connection is closed. floodConns := make([]net.Conn, maxPending) for i := 0; i < maxPending; i++ { - c, err := net.DialTimeout("tcp", hostAddr, 2*time.Second) + dialer := net.Dialer{Timeout: 2 * time.Second, + LocalAddr: &net.TCPAddr{IP: net.ParseIP(fmt.Sprintf("127.0.0.%d", i+2))}} + c, err := dialer.Dial("tcp", hostAddr) if err != nil { t.Fatalf("flood conn %d: dial: %v", i, err) } @@ -2488,7 +2514,9 @@ func TestConcurrentHandshakeFlood(t *testing.T) { defer wg.Done() readyCh <- struct{}{} // signal: I am at the gate <-startCh // wait for the simultaneous release - c, dialErr := net.DialTimeout("tcp", hostAddr, 3*time.Second) + dialer := net.Dialer{Timeout: 3 * time.Second, + LocalAddr: &net.TCPAddr{IP: net.ParseIP(fmt.Sprintf("127.0.0.%d", idx+2))}} + c, dialErr := dialer.Dial("tcp", hostAddr) if dialErr != nil { if mustSucceed { // Non-fatal from here; caller checks the slice. diff --git a/p2p/handshake_budget.go b/p2p/handshake_budget.go new file mode 100644 index 0000000..3e837f5 --- /dev/null +++ b/p2p/handshake_budget.go @@ -0,0 +1,47 @@ +package p2p + +// Reservations cover all inbound pre-registration handshake work. A single +// source cannot consume the global budget; established-peer checks remain +// independently enforced at registration. +func (h *Host) reserveInboundHandshake(ip string) bool { + if ip == "" { + return false + } + h.handshakeMu.Lock() + defer h.handshakeMu.Unlock() + limit := 3 + if h.cfg.MaxPeersPerIP > 0 && h.cfg.MaxPeersPerIP < limit { + limit = h.cfg.MaxPeersPerIP + } + if total := h.cfg.MaxPendingHandshakes; total > 1 && limit > total/2 { + limit = total / 2 + } + if h.pendingHandshakeIPs[ip] >= limit { + return false + } + if h.cfg.MaxPendingHandshakes > 0 && + h.pendingHandshakes.Load() >= int64(h.cfg.MaxPendingHandshakes) { + return false + } + if h.pendingHandshakeIPs == nil { + h.pendingHandshakeIPs = make(map[string]int) + } + h.pendingHandshakeIPs[ip]++ + h.pendingHandshakes.Add(1) + return true +} + +func (h *Host) releaseInboundHandshake(ip string) { + h.handshakeMu.Lock() + defer h.handshakeMu.Unlock() + n := h.pendingHandshakeIPs[ip] + if n == 0 { + return + } + if n == 1 { + delete(h.pendingHandshakeIPs, ip) + } else { + h.pendingHandshakeIPs[ip] = n - 1 + } + h.pendingHandshakes.Add(-1) +} diff --git a/p2p/handshake_budget_test.go b/p2p/handshake_budget_test.go new file mode 100644 index 0000000..0e1bee7 --- /dev/null +++ b/p2p/handshake_budget_test.go @@ -0,0 +1,60 @@ +package p2p + +import ( + "sync" + "testing" +) + +func TestInboundHandshakeBudgetBySource(t *testing.T) { + h := &Host{cfg: Config{MaxPendingHandshakes: 20, MaxPeersPerIP: 3}} + for i := 0; i < 3; i++ { + if !h.reserveInboundHandshake("one") { + t.Fatal("available source reservation denied") + } + } + for i := 0; i < 30; i++ { + if h.reserveInboundHandshake("one") { + t.Fatal("source quota exceeded") + } + } + if !h.reserveInboundHandshake("two") { + t.Fatal("independent source denied") + } + h.releaseInboundHandshake("two") + for i := 0; i < 3; i++ { + h.releaseInboundHandshake("one") + } + h.releaseInboundHandshake("one") + if h.PendingHandshakes() != 0 || len(h.pendingHandshakeIPs) != 0 { + t.Fatal("reservation state retained after release") + } +} + +func TestInboundHandshakeBudgetConcurrent(t *testing.T) { + h := &Host{cfg: Config{MaxPendingHandshakes: 4}} + var wg sync.WaitGroup + for i := 0; i < 50; i++ { + wg.Add(1) + go func() { + defer wg.Done() + h.reserveInboundHandshake("source") + }() + } + wg.Wait() + if h.PendingHandshakes() != 2 { + t.Fatalf("one source used %d reservations", h.PendingHandshakes()) + } + if !h.reserveInboundHandshake("other") || !h.reserveInboundHandshake("other") { + t.Fatal("reserved capacity unavailable to independent source") + } + if h.reserveInboundHandshake("third") { + t.Fatal("global quota exceeded") + } + for _, ip := range []string{"source", "other"} { + h.releaseInboundHandshake(ip) + h.releaseInboundHandshake(ip) + } + if h.PendingHandshakes() != 0 { + t.Fatal("global reservation leak") + } +} diff --git a/p2p/host.go b/p2p/host.go index 7111b20..4d47590 100644 --- a/p2p/host.go +++ b/p2p/host.go @@ -210,6 +210,16 @@ type Handler interface { GetBlock(hash crypto.Hash32) *core.Block } +// BlockValidationHandler validates gossip without changing asynchronous +// consensus admission or the wire protocol. +type BlockValidationHandler interface { + ValidateBlock(*core.Block) error +} + +type SyncHeaderValidationHandler interface { + ValidateSyncHeader(core.BlockHeader) error +} + // pendingBlockEntry records one outstanding MsgGetBlock request so the // stall-detection ticker can log actionable diagnostics if no MsgBlock arrives. type pendingBlockEntry struct { @@ -507,7 +517,9 @@ type Host struct { // executing the TLS handshake. Guarded by MaxPendingHandshakes; // uses atomic ops so acceptLoop and handleConn coordinate without // holding h.mu. - pendingHandshakes atomic.Int64 + pendingHandshakes atomic.Int64 + handshakeMu sync.Mutex + pendingHandshakeIPs map[string]int // badBlockCounts tracks out-of-range-block strikes per remote IP. // Entries expire after badBlockStrikeTTL and the map is capped at @@ -1994,10 +2006,8 @@ func (h *Host) acceptLoop() { // hold one goroutine per connection for up to 10 s. // MaxPendingHandshakes caps the total in-flight handshakes so // the node cannot be goroutine-starved by a connect-flood. - if h.cfg.MaxPendingHandshakes > 0 && h.cfg.TLSConfig != nil { - cur := h.pendingHandshakes.Add(1) - if cur > int64(h.cfg.MaxPendingHandshakes) { - h.pendingHandshakes.Add(-1) + if h.cfg.TLSConfig != nil { + if !h.reserveInboundHandshake(connIP(conn.RemoteAddr().String())) { h.log.Info("MaxPendingHandshakes reached — inbound connection rejected", "limit", h.cfg.MaxPendingHandshakes) conn.Close() @@ -2652,14 +2662,16 @@ func (h *Host) handleConn(conn net.Conn, outbound bool, dialID, pendingID uint64 // below). releaseHS is idempotent; calling it more than once is safe. // We also call it explicitly right after a successful handshake so that // the slot is freed as early as possible rather than at connection close. - hsSlotHeld := !outbound && h.cfg.MaxPendingHandshakes > 0 && h.cfg.TLSConfig != nil + hsSlotHeld := !outbound && h.cfg.TLSConfig != nil releaseHS := func() { if hsSlotHeld { - h.pendingHandshakes.Add(-1) + h.releaseInboundHandshake(connIP(conn.RemoteAddr().String())) hsSlotHeld = false } } defer releaseHS() // safety net: covers ban-check return and any other early exit + handshakeTimer := time.AfterFunc(h.cfg.HandshakeTimeout, func() { conn.Close() }) + defer handshakeTimer.Stop() // Reject banned peers immediately if h.mgr.IsBanned(addr) { @@ -2683,7 +2695,6 @@ func (h *Host) handleConn(conn net.Conn, outbound bool, dialID, pendingID uint64 conn.Close() return } - releaseHS() // handshake complete — free the slot early, before message loop tlsConn.SetDeadline(time.Time{}) //nolint:errcheck fp := PeerFingerprint(conn) h.log.Debug("tls handshake ok", "addr", addr, "fingerprint", fp) @@ -2894,6 +2905,9 @@ func (h *Host) handleConn(conn net.Conn, outbound bool, dialID, pendingID uint64 "peer_height", peerHeight, "direction", map[bool]string{true: "out", false: "in"}[outbound], ) + // Keep inbound reservations through TLS, identity and application handshake. + handshakeTimer.Stop() + releaseHS() // Initiate header sync only when the peer is not known to be behind us. // Requesting our tip locator from a lagging peer makes its unknown-locator @@ -3189,6 +3203,9 @@ func (h *Host) dispatch(peer *Peer, msgType MessageType, data []byte) error { if block != nil { ourTip := h.handler.CurrentHeight() peerIP := connIP(peer.addr) + blockValidated := false + var validationErr error + trustedSyncHeader := false // A syncing peer may gossip canonical blocks back to the node that // supplied them, especially when the same validators have parallel @@ -3316,19 +3333,18 @@ func (h *Host) dispatch(peer *Peer, msgType MessageType, data []byte) error { // bare IP so a reconnect on a new source port does not bypass the // enforcement. // - // A strike is only warranted when the received block is far ahead of - // our tip AND the sending peer itself is NOT far ahead of us. A peer - // that is genuinely further along (peer.height > ourTip+lead) is a - // node we are actively syncing from; the blocks it sends at its own - // tip arrive via gossip before our sync pipeline has applied the - // intermediate blocks. Counting those as strikes would permanently - // ban the validator we are catching up to. - // - // Rogue peers that fabricate future-height blocks announce a - // peer.height at or below our tip (they pretend to be at the same - // height), so the second condition catches them. - if liveHeightLead := h.badBlockHeightLeadV.Load(); block.Header.Height > ourTip+liveHeightLead && - peer.height <= ourTip+liveHeightLead { + // Fully verified data may proceed; remote height is informational only. + // Whitelist policy remains independent from verification and relay policy. + // Locally verified sync data, not advertised height, grants an exemption. + if validator, ok := h.handler.(BlockValidationHandler); ok { + validationErr = validator.ValidateBlock(block) + blockValidated = validationErr == nil + } + if validator, ok := h.handler.(SyncHeaderValidationHandler); ok { + trustedSyncHeader = validator.ValidateSyncHeader(block.Header) == nil + } + if liveHeightLead := h.badBlockHeightLeadV.Load(); block.Header.Height > ourTip && + block.Header.Height-ourTip > liveHeightLead && !blockValidated && !trustedSyncHeader { // Whitelisted peers are trusted validators; skip the // strike counter entirely so a temporarily-ahead validator // is never auto-banned for being on a longer fork. @@ -3437,6 +3453,17 @@ func (h *Host) dispatch(peer *Peer, msgType MessageType, data []byte) error { h.badBlockMu.Unlock() processBlock: + // Also validate branches that skipped directly here under whitelist policy. + if !blockValidated && validationErr == nil { + if validator, ok := h.handler.(BlockValidationHandler); ok { + validationErr = validator.ValidateBlock(block) + blockValidated = validationErr == nil + } + } + if validationErr != nil { + h.log.Debug("block admission rejected", "err", validationErr) + return nil + } // Rate-limit block ingestion per peer. The token bucket // sleeps the dispatch goroutine (which is the only reader // for this peer's conn) when tokens are exhausted, creating @@ -3479,7 +3506,7 @@ func (h *Host) dispatch(peer *Peer, msgType MessageType, data []byte) error { h.requestHeaders(peer) } } - if isNew { + if isNew && blockValidated { sb := blockToMsg(block) fromAddr := peer.addr h.mu.RLock()