diff --git a/cmd/geth/main.go b/cmd/geth/main.go index e099755..819572b 100644 --- a/cmd/geth/main.go +++ b/cmd/geth/main.go @@ -169,6 +169,10 @@ var ( utils.MinerNotifyFullFlag, utils.ECBP1100Flag, utils.ECBP1100NoDisableFlag, + utils.SCDOCheckpointSignersFlag, + utils.SCDOCheckpointThresholdFlag, + utils.SCDOCheckpointURLFlag, + utils.SCDOCheckpointPollFlag, utils.OverrideECBP1100DeactivateFlag, configFileFlag, utils.LogDebugFlag, diff --git a/cmd/scdocheckpoint/main.go b/cmd/scdocheckpoint/main.go new file mode 100644 index 0000000..42a2ebc --- /dev/null +++ b/cmd/scdocheckpoint/main.go @@ -0,0 +1,218 @@ +// scdocheckpoint: SCDO shard0 checkpoint signer / verifier (9Y9 PTY LTD). +// +// scdocheckpoint genkey -out signer.key (throwaway/test key, prints address) +// scdocheckpoint sign -key signer.key -rpc URL[,URL] -every 100 -depth 12 -out cp.json [-push URL,...] [-loop 5s] +// scdocheckpoint signraw -key k -chainid N -number N -hash 0x.. -out cp.json (tests: forge/tamper) +// scdocheckpoint verify -file cp.json -signers 0x..,0x.. -threshold 1 -chainid 5680 +// +// "sign" only signs a block that every -rpc node reports with the same hash at +// (head - depth) rounded down to a multiple of -every. +package main + +import ( + "bytes" + "encoding/json" + "flag" + "fmt" + "math/big" + "net/http" + "os" + "path/filepath" + "strings" + "time" + + "github.com/ethereum/go-ethereum/common" + "github.com/ethereum/go-ethereum/common/hexutil" + "github.com/ethereum/go-ethereum/core" + "github.com/ethereum/go-ethereum/crypto" +) + +func die(f string, a ...interface{}) { fmt.Fprintf(os.Stderr, f+"\n", a...); os.Exit(1) } + +func rpc(url, method string, params []interface{}, out interface{}) error { + if params == nil { + params = []interface{}{} + } + body, _ := json.Marshal(map[string]interface{}{"jsonrpc": "2.0", "id": 1, "method": method, "params": params}) + resp, err := (&http.Client{Timeout: 10 * time.Second}).Post(url, "application/json", bytes.NewReader(body)) + if err != nil { + return err + } + defer resp.Body.Close() + var r struct { + Result json.RawMessage `json:"result"` + Error *struct{ Message string } `json:"error"` + } + if err := json.NewDecoder(resp.Body).Decode(&r); err != nil { + return err + } + if r.Error != nil { + return fmt.Errorf("%s", r.Error.Message) + } + return json.Unmarshal(r.Result, out) +} + +func writeAtomic(path string, v interface{}) { + b, _ := json.MarshalIndent(v, "", " ") + tmp := path + ".tmp" + if err := os.WriteFile(tmp, b, 0o644); err != nil { + die("write: %v", err) + } + if err := os.Rename(tmp, path); err != nil { + die("rename: %v", err) + } +} + +func signCP(keyfile string, cp *core.SCDOCheckpoint) { + key, err := crypto.LoadECDSA(keyfile) + if err != nil { + die("load key: %v", err) + } + cp.Version = 1 + cp.Timestamp = uint64(time.Now().Unix()) + cp.Message = cp.SigningMessage() + if err := core.SignSCDOCheckpoint(cp, func(h []byte) ([]byte, error) { return crypto.Sign(h, key) }); err != nil { + die("sign: %v", err) + } + cp.Signers = []common.Address{crypto.PubkeyToAddress(key.PublicKey)} +} + +func main() { + if len(os.Args) < 2 { + die("usage: scdocheckpoint genkey|sign|signraw|verify ...") + } + fs := flag.NewFlagSet(os.Args[1], flag.ExitOnError) + out := fs.String("out", "", "output file") + keyf := fs.String("key", "", "signer key file (hex)") + rpcs := fs.String("rpc", "http://127.0.0.1:8545", "comma-separated trusted node RPC URLs (all must agree)") + every := fs.Uint64("every", 100, "checkpoint interval (blocks)") + depth := fs.Uint64("depth", 12, "confirmation depth before signing") + push := fs.String("push", "", "comma-separated node RPC URLs to push scdo_submitCheckpoint to") + loop := fs.Duration("loop", 0, "repeat every interval (0 = once)") + chainid := fs.Uint64("chainid", 0, "chain id") + number := fs.Uint64("number", 0, "block number") + hash := fs.String("hash", "", "block hash") + file := fs.String("file", "", "checkpoint file") + signers := fs.String("signers", "", "authorised signer addresses") + threshold := fs.Int("threshold", 1, "signature threshold") + serve := fs.String("serve", "", "sign: also serve the -out directory read-only over HTTP on this address (e.g. 0.0.0.0:18590)") + fs.Parse(os.Args[2:]) + + switch os.Args[1] { + case "genkey": + key, _ := crypto.GenerateKey() + if *out == "" { + die("-out required") + } + if err := crypto.SaveECDSA(*out, key); err != nil { + die("save: %v", err) + } + os.Chmod(*out, 0o600) + fmt.Println(crypto.PubkeyToAddress(key.PublicKey).Hex()) + case "signraw": + cp := &core.SCDOCheckpoint{ChainID: *chainid, Number: *number, Hash: common.HexToHash(*hash)} + signCP(*keyf, cp) + writeAtomic(*out, cp) + case "verify": + b, err := os.ReadFile(*file) + if err != nil { + die("%v", err) + } + var cp core.SCDOCheckpoint + if err := json.Unmarshal(b, &cp); err != nil { + die("%v", err) + } + pol := &core.SCDOCheckpointPolicy{ChainID: *chainid, Threshold: *threshold} + for _, a := range strings.Split(*signers, ",") { + pol.Signers = append(pol.Signers, common.HexToAddress(strings.TrimSpace(a))) + } + who, err := pol.Verify(&cp) + if err != nil { + die("INVALID: %v", err) + } + fmt.Printf("VALID number=%d hash=%s signers=%v\n", cp.Number, cp.Hash.Hex(), who) + case "sign": + urls := strings.Split(*rpcs, ",") + if *serve != "" { + dir := filepath.Dir(*out) + fsrv := http.FileServer(http.Dir(dir)) + mux := http.NewServeMux() + mux.HandleFunc("/", func(w http.ResponseWriter, r *http.Request) { + if r.Method != http.MethodGet && r.Method != http.MethodHead { + http.Error(w, "method not allowed", http.StatusMethodNotAllowed) + return + } + if !strings.HasSuffix(r.URL.Path, ".json") { + http.NotFound(w, r) + return + } + w.Header().Set("Cache-Control", "no-cache") + w.Header().Set("Access-Control-Allow-Origin", "*") + fsrv.ServeHTTP(w, r) + }) + srv := &http.Server{Addr: *serve, Handler: mux, ReadTimeout: 10 * time.Second, WriteTimeout: 10 * time.Second} + go func() { die("serve: %v", srv.ListenAndServe()) }() + } + var last uint64 + if b, err := os.ReadFile(*out); err == nil { + var prev core.SCDOCheckpoint + if json.Unmarshal(b, &prev) == nil { + last = prev.Number + } + } + for { + var head hexutil.Uint64 + var cid hexutil.Big + err := rpc(urls[0], "eth_blockNumber", nil, &head) + if err == nil { + err = rpc(urls[0], "eth_chainId", nil, &cid) + } + if err != nil { + fmt.Fprintln(os.Stderr, "rpc:", err) + } else if uint64(head) > *depth { + target := (uint64(head) - *depth) / *every * *every + if target > last && target > 0 { + var agreed common.Hash + ok := true + for _, u := range urls { + var blk struct{ Hash common.Hash } + if err := rpc(u, "eth_getBlockByNumber", []interface{}{hexutil.EncodeUint64(target), false}, &blk); err != nil { + fmt.Fprintln(os.Stderr, "rpc:", u, err) + ok = false + break + } + if agreed == (common.Hash{}) { + agreed = blk.Hash + } else if agreed != blk.Hash { + fmt.Fprintf(os.Stderr, "nodes disagree at %d: %s vs %s; not signing\n", target, agreed.Hex(), blk.Hash.Hex()) + ok = false + break + } + } + if ok { + cp := &core.SCDOCheckpoint{ChainID: (*big.Int)(&cid).Uint64(), Number: target, Hash: agreed} + signCP(*keyf, cp) + writeAtomic(*out, cp) + writeAtomic(filepath.Join(filepath.Dir(*out), fmt.Sprintf("checkpoint-%d.json", target)), cp) + last = target + fmt.Printf("%s signed checkpoint number=%d hash=%s\n", time.Now().Format(time.RFC3339), target, agreed.Hex()) + for _, p := range strings.Split(*push, ",") { + if p = strings.TrimSpace(p); p != "" { + var res map[string]interface{} + if err := rpc(p, "scdo_submitCheckpoint", []interface{}{cp}, &res); err != nil { + fmt.Fprintln(os.Stderr, "push:", p, err) + } + } + } + } + } + } + if *loop == 0 { + return + } + time.Sleep(*loop) + } + default: + die("unknown command %s", os.Args[1]) + } +} diff --git a/cmd/utils/flags.go b/cmd/utils/flags.go index 7d24df6..bda0d93 100644 --- a/cmd/utils/flags.go +++ b/cmd/utils/flags.go @@ -1085,6 +1085,28 @@ Please note that --` + MetricsHTTPFlag.Name + ` must be set to start the server. Usage: "Manually specify the ECBP-1100 (MESS) deactivation block number, overriding the bundled setting", Category: flags.EthCategory, } + SCDOCheckpointSignersFlag = &cli.StringFlag{ + Name: "scdo.checkpoint.signers", + Usage: "Comma-separated addresses of authorised SCDO checkpoint signers (enables signed checkpoints)", + Category: flags.EthCategory, + } + SCDOCheckpointThresholdFlag = &cli.IntFlag{ + Name: "scdo.checkpoint.threshold", + Usage: "Number of distinct authorised signatures a checkpoint needs", + Value: 1, + Category: flags.EthCategory, + } + SCDOCheckpointURLFlag = &cli.StringFlag{ + Name: "scdo.checkpoint.url", + Usage: "Comma-separated checkpoint sources (https://..., http://..., file:///...) polled for the latest signed checkpoint", + Category: flags.EthCategory, + } + SCDOCheckpointPollFlag = &cli.DurationFlag{ + Name: "scdo.checkpoint.poll", + Usage: "Checkpoint source poll interval", + Value: 30 * time.Second, + Category: flags.EthCategory, + } ECBP1100NoDisableFlag = &cli.BoolFlag{ Name: "ecbp1100.nodisable", Usage: "Short-circuit ECBP-1100 (MESS) disable mechanisms; (yields a permanent-once-activated state, deactivating auto-shutoff mechanisms)", @@ -1956,6 +1978,7 @@ func SetEthConfig(ctx *cli.Context, stack *node.Node, cfg *ethconfig.Config) { setMiner(ctx, &cfg.Miner) setRequiredBlocks(ctx, cfg) setLes(ctx, cfg) + setSCDOCheckpoints(ctx, cfg) // Cap the cache allowance and tune the garbage collector mem, err := gopsutil.VirtualMemory() @@ -2641,3 +2664,25 @@ func MakeTrieDatabase(ctx *cli.Context, disk ethdb.Database, preimage bool, read } return triedb.NewDatabase(disk, config) } + +// setSCDOCheckpoints applies the SCDO signed checkpoint flags. +func setSCDOCheckpoints(ctx *cli.Context, cfg *ethconfig.Config) { + if v := ctx.String(SCDOCheckpointSignersFlag.Name); v != "" { + cfg.SCDOCheckpointSigners = nil + for _, a := range strings.Split(v, ",") { + a = strings.TrimSpace(a) + if a == "" { + continue + } + if !common.IsHexAddress(a) { + Fatalf("Invalid --%s address: %q", SCDOCheckpointSignersFlag.Name, a) + } + cfg.SCDOCheckpointSigners = append(cfg.SCDOCheckpointSigners, common.HexToAddress(a)) + } + } + cfg.SCDOCheckpointThreshold = ctx.Int(SCDOCheckpointThresholdFlag.Name) + if v := ctx.String(SCDOCheckpointURLFlag.Name); v != "" { + cfg.SCDOCheckpointURLs = strings.Split(v, ",") + } + cfg.SCDOCheckpointInterval = ctx.Duration(SCDOCheckpointPollFlag.Name) +} diff --git a/core/block_validator.go b/core/block_validator.go index bedb81f..b94f3b7 100644 --- a/core/block_validator.go +++ b/core/block_validator.go @@ -56,6 +56,10 @@ func (v *BlockValidator) ValidateBody(block *types.Block) error { if v.bc.HasBlockAndState(block.Hash(), block.NumberU64()) { return ErrKnownBlock } + // SCDO signed checkpoints: never import a block conflicting with a checkpoint. + if err := v.bc.scdoCheckBlock(block.Header()); err != nil { + return err + } // Header validity is known at this point. Here we verify that uncles, transactions // and withdrawals given in the block body match the header. diff --git a/core/blockchain.go b/core/blockchain.go index 7b271d5..7e3c48c 100644 --- a/core/blockchain.go +++ b/core/blockchain.go @@ -264,6 +264,8 @@ type BlockChain struct { artificialFinalityNoDisable *int32 // manual override prevents disabling artificial finality feature activation artificialFinalityEnabledStatus int32 // toggles artificial finality features; will be always 1 if artificialFinalityForce=1 + + scdoCP atomic.Pointer[scdoCheckpointer] // SCDO signed checkpoints (nil = disabled) } // NewBlockChain returns a fully initialised block chain using information diff --git a/core/forkchoice.go b/core/forkchoice.go index d07cdf6..6258869 100644 --- a/core/forkchoice.go +++ b/core/forkchoice.go @@ -149,6 +149,13 @@ func (f *ForkChoice) ReorgNeeded(current *types.Header, extern *types.Header) (b } } + // SCDO signed checkpoints: a validly signed checkpoint overrides TD. + if bc, ok := f.chain.(*BlockChain); ok { + if decided, cpReorg := bc.scdoForkChoice(current, extern); decided { + return cpReorg, nil + } + } + // If reorg is not needed (false), then we can just return. // The following logic adds a condition only in the case where a reorg would // otherwise be indicated. diff --git a/core/scdo_checkpoint.go b/core/scdo_checkpoint.go new file mode 100644 index 0000000..0e3577a --- /dev/null +++ b/core/scdo_checkpoint.go @@ -0,0 +1,407 @@ +// Copyright 2026 9Y9 PTY LTD (SCDO shard0). LGPL-3.0, same as go-ethereum. +// +// SCDO signed checkpoints ("company finality" overlay for Ethash PoW). +// +// A configured set of checkpoint signer keys signs (chainId, number, hash) of +// blocks that are already buried under the signer node's head (e.g. every 100 +// blocks at depth 12). A node that holds a valid signed checkpoint: +// +// 1. refuses to import any block at the checkpoint height with another hash +// (ValidateBody -> ErrSCDOCheckpointMismatch), so a conflicting fork can +// never be extended past the checkpoint; +// 2. refuses any reorg that would drop the checkpointed block from the +// canonical chain (ForkChoice.ReorgNeeded), regardless of total difficulty; +// 3. prefers a chain that contains the checkpoint over one that conflicts with it; +// 4. persists the latest checkpoint in the chain database (survives restarts) +// and marks the checkpoint block as the "finalized" block (eth_getBlockByNumber("finalized")). +// +// Mining, block validity rules and the p2p protocol are unchanged. If no +// checkpoint arrives (signer offline), the node behaves exactly like upstream. +// +// Signature scheme: EIP-191 personal_sign over the ASCII message +// "SCDO-CHECKPOINT-v1 chainId= number= hash=<0x..>" +// so any Ethereum wallet / HSM / ethers.verifyMessage can produce and check it. + +package core + +import ( + "encoding/json" + "errors" + "fmt" + "math/big" + "sort" + "strings" + "sync" + "sync/atomic" + + "github.com/ethereum/go-ethereum/common" + "github.com/ethereum/go-ethereum/common/hexutil" + "github.com/ethereum/go-ethereum/core/types" + "github.com/ethereum/go-ethereum/crypto" + "github.com/ethereum/go-ethereum/log" +) + +var ( + ErrSCDOCheckpointMismatch = errors.New("block conflicts with SCDO signed checkpoint") + ErrSCDOCheckpointVersion = errors.New("unsupported checkpoint version") + ErrSCDOCheckpointChainID = errors.New("checkpoint chainId mismatch") + ErrSCDOCheckpointNoSig = errors.New("checkpoint has no signatures") + ErrSCDOCheckpointBadSig = errors.New("invalid checkpoint signature") + ErrSCDOCheckpointSigner = errors.New("checkpoint signed by unknown signer") + ErrSCDOCheckpointThreshold = errors.New("not enough distinct authorised checkpoint signatures") + ErrSCDOCheckpointOld = errors.New("checkpoint not newer than the current one") + ErrSCDOCheckpointDisabled = errors.New("SCDO checkpoints not enabled on this node") + ErrSCDOCheckpointZeroHash = errors.New("checkpoint hash is zero") + scdoCheckpointDBKey = []byte("scdo-signed-checkpoint-latest") + scdoCheckpointMessagePrefix = "SCDO-CHECKPOINT-v1" +) + +// SCDOCheckpoint is a signed (number, hash) pair. JSON is the wire/file format. +type SCDOCheckpoint struct { + Version int `json:"version"` + ChainID uint64 `json:"chainId"` + Number uint64 `json:"number"` + Hash common.Hash `json:"hash"` + Signatures []hexutil.Bytes `json:"signatures"` + Signers []common.Address `json:"signers,omitempty"` // informational only, never trusted + Timestamp uint64 `json:"timestamp,omitempty"` // informational: unix time of signing + Message string `json:"message,omitempty"` // informational: the signed text +} + +// SigningMessage returns the exact text that is signed (EIP-191 personal_sign). +func (c *SCDOCheckpoint) SigningMessage() string { + return fmt.Sprintf("%s chainId=%d number=%d hash=%s", scdoCheckpointMessagePrefix, c.ChainID, c.Number, c.Hash.Hex()) +} + +// SigningHash is keccak256("\x19Ethereum Signed Message:\n" + len(msg) + msg). +func (c *SCDOCheckpoint) SigningHash() []byte { + msg := c.SigningMessage() + return crypto.Keccak256([]byte(fmt.Sprintf("\x19Ethereum Signed Message:\n%d%s", len(msg), msg))) +} + +// SignSCDOCheckpoint appends a signature made with a raw secp256k1 key (tooling/tests). +func SignSCDOCheckpoint(c *SCDOCheckpoint, sign func(hash []byte) ([]byte, error)) error { + sig, err := sign(c.SigningHash()) + if err != nil { + return err + } + if len(sig) != 65 { + return ErrSCDOCheckpointBadSig + } + if sig[64] < 27 { + sig[64] += 27 // personal_sign convention + } + c.Signatures = append(c.Signatures, sig) + return nil +} + +// SCDOCheckpointPolicy defines who may sign checkpoints for which chain. +type SCDOCheckpointPolicy struct { + ChainID uint64 + Signers []common.Address + Threshold int +} + +func (p *SCDOCheckpointPolicy) String() string { + s := make([]string, len(p.Signers)) + for i, a := range p.Signers { + s[i] = a.Hex() + } + return fmt.Sprintf("chainId=%d threshold=%d/%d signers=[%s]", p.ChainID, p.Threshold, len(p.Signers), strings.Join(s, ",")) +} + +// Verify checks a checkpoint against the policy and returns the distinct authorised signers. +func (p *SCDOCheckpointPolicy) Verify(c *SCDOCheckpoint) ([]common.Address, error) { + if c == nil { + return nil, errors.New("nil checkpoint") + } + if c.Version != 1 { + return nil, ErrSCDOCheckpointVersion + } + if c.ChainID != p.ChainID { + return nil, fmt.Errorf("%w: have %d want %d", ErrSCDOCheckpointChainID, c.ChainID, p.ChainID) + } + if c.Hash == (common.Hash{}) { + return nil, ErrSCDOCheckpointZeroHash + } + if len(c.Signatures) == 0 { + return nil, ErrSCDOCheckpointNoSig + } + if len(c.Signatures) > 32 { + return nil, fmt.Errorf("%w: too many signatures", ErrSCDOCheckpointBadSig) + } + allowed := make(map[common.Address]bool, len(p.Signers)) + for _, a := range p.Signers { + allowed[a] = true + } + hash := c.SigningHash() + seen := make(map[common.Address]bool) + var good []common.Address + for i, sig := range c.Signatures { + if len(sig) != 65 { + return nil, fmt.Errorf("%w: signature %d has length %d", ErrSCDOCheckpointBadSig, i, len(sig)) + } + s := make([]byte, 65) + copy(s, sig) + if s[64] >= 27 { + s[64] -= 27 + } + if s[64] > 1 { + return nil, fmt.Errorf("%w: signature %d bad v", ErrSCDOCheckpointBadSig, i) + } + // Reject malleable (high-s) signatures. + if !crypto.ValidateSignatureValues(s[64], new(big.Int).SetBytes(s[0:32]), new(big.Int).SetBytes(s[32:64]), true) { + return nil, fmt.Errorf("%w: signature %d invalid values", ErrSCDOCheckpointBadSig, i) + } + pub, err := crypto.SigToPub(hash, s) + if err != nil { + return nil, fmt.Errorf("%w: signature %d: %v", ErrSCDOCheckpointBadSig, i, err) + } + addr := crypto.PubkeyToAddress(*pub) + if !allowed[addr] { + return nil, fmt.Errorf("%w: %s", ErrSCDOCheckpointSigner, addr.Hex()) + } + if !seen[addr] { + seen[addr] = true + good = append(good, addr) + } + } + if len(good) < p.Threshold { + return nil, fmt.Errorf("%w: have %d need %d", ErrSCDOCheckpointThreshold, len(good), p.Threshold) + } + sort.Slice(good, func(i, j int) bool { return good[i].Hex() < good[j].Hex() }) + return good, nil +} + +// scdoCheckpointer is the per-BlockChain checkpoint state. +type scdoCheckpointer struct { + mu sync.Mutex + policy *SCDOCheckpointPolicy + latest atomic.Pointer[SCDOCheckpoint] + rejected atomic.Uint64 // blocks/reorgs refused because of a checkpoint + conflict atomic.Bool // local canonical chain conflicts with the latest checkpoint +} + +// EnableSCDOCheckpoints activates the checkpoint rules with the given policy +// and restores the persisted checkpoint (re-verified against the policy). +func (bc *BlockChain) EnableSCDOCheckpoints(policy *SCDOCheckpointPolicy) error { + if policy == nil || len(policy.Signers) == 0 { + return errors.New("SCDO checkpoints: empty signer list") + } + if policy.Threshold < 1 || policy.Threshold > len(policy.Signers) { + return fmt.Errorf("SCDO checkpoints: bad threshold %d for %d signers", policy.Threshold, len(policy.Signers)) + } + cp := &scdoCheckpointer{policy: policy} + if blob, err := bc.db.Get(scdoCheckpointDBKey); err == nil && len(blob) > 0 { + var stored SCDOCheckpoint + if err := json.Unmarshal(blob, &stored); err != nil { + log.Warn("SCDO checkpoint: stored checkpoint unreadable, ignoring", "err", err) + } else if _, err := policy.Verify(&stored); err != nil { + log.Warn("SCDO checkpoint: stored checkpoint no longer valid for signer policy, ignoring", "number", stored.Number, "err", err) + } else { + cp.latest.Store(&stored) + log.Info("SCDO checkpoint restored from database", "number", stored.Number, "hash", stored.Hash) + } + } + bc.scdoCP.Store(cp) + log.Info("SCDO signed checkpoints enabled", "policy", policy.String()) + bc.SCDOCheckpointMaintain() + return nil +} + +func (bc *BlockChain) scdo() *scdoCheckpointer { return bc.scdoCP.Load() } + +// SCDOCheckpoint returns the latest accepted checkpoint (nil if none / disabled). +func (bc *BlockChain) SCDOCheckpoint() *SCDOCheckpoint { + if s := bc.scdo(); s != nil { + return s.latest.Load() + } + return nil +} + +// SCDOCheckpointPolicy returns the active policy (nil if disabled). +func (bc *BlockChain) SCDOCheckpointPolicy() *SCDOCheckpointPolicy { + if s := bc.scdo(); s != nil { + return s.policy + } + return nil +} + +// SCDOCheckpointStats returns (#rejected blocks/reorgs, local chain conflicts with checkpoint). +func (bc *BlockChain) SCDOCheckpointStats() (uint64, bool) { + if s := bc.scdo(); s != nil { + return s.rejected.Load(), s.conflict.Load() + } + return 0, false +} + +// AddSCDOCheckpoint verifies and, if newer, adopts and persists a checkpoint. +// Returns (true, nil) when adopted, (false, nil) when it equals the current one. +func (bc *BlockChain) AddSCDOCheckpoint(c *SCDOCheckpoint) (bool, error) { + s := bc.scdo() + if s == nil { + return false, ErrSCDOCheckpointDisabled + } + signers, err := s.policy.Verify(c) + if err != nil { + return false, err + } + s.mu.Lock() + cur := s.latest.Load() + if cur != nil { + if c.Number == cur.Number && c.Hash == cur.Hash { + s.mu.Unlock() + return false, nil + } + if c.Number <= cur.Number { + s.mu.Unlock() + if c.Number == cur.Number { + log.Error("SCDO checkpoint: validly signed checkpoint CONFLICTS with current one (signer key compromise?)", "number", c.Number, "have", cur.Hash, "got", c.Hash) + } + return false, fmt.Errorf("%w: have %d got %d", ErrSCDOCheckpointOld, cur.Number, c.Number) + } + // A new checkpoint must descend from the previous one when both are known locally. + if h := bc.GetHeaderByHash(c.Hash); h != nil && h.Number.Uint64() == c.Number { + if st := bc.scdoChainHas(h, cur); st == scdoCPConflict { + s.mu.Unlock() + log.Error("SCDO checkpoint: new checkpoint does not descend from previous checkpoint, refusing", "number", c.Number, "hash", c.Hash, "prev", cur.Number) + return false, errors.New("checkpoint does not descend from previous checkpoint") + } + } + } + stored := *c + stored.Signers = signers + stored.Message = c.SigningMessage() + blob, _ := json.Marshal(&stored) + if err := bc.db.Put(scdoCheckpointDBKey, blob); err != nil { + s.mu.Unlock() + return false, err + } + s.latest.Store(&stored) + s.mu.Unlock() + log.Info("SCDO signed checkpoint accepted", "number", c.Number, "hash", c.Hash, "signers", len(signers)) + bc.SCDOCheckpointMaintain() + return true, nil +} + +// SCDOCheckpointMaintain updates the finalized marker and detects / heals a +// local canonical chain that conflicts with the latest checkpoint. +func (bc *BlockChain) SCDOCheckpointMaintain() { + s := bc.scdo() + if s == nil { + return + } + cp := s.latest.Load() + if cp == nil { + return + } + head := bc.CurrentBlock() + if head == nil || head.Number.Uint64() < cp.Number { + s.conflict.Store(false) + return + } + canon := bc.GetCanonicalHash(cp.Number) + if canon == cp.Hash { + s.conflict.Store(false) + if fin := bc.CurrentFinalBlock(); fin == nil || fin.Hash() != cp.Hash { + if h := bc.GetHeaderByHash(cp.Hash); h != nil { + bc.SetFinalized(h) + } + } + return + } + // Canonical chain conflicts with a validly signed checkpoint. + s.conflict.Store(true) + block := bc.GetBlockByHash(cp.Hash) + if block == nil || !bc.HasState(block.Root()) { + log.Error("SCDO checkpoint: local chain CONFLICTS with signed checkpoint and checkpoint block is unknown; waiting for peers on the checkpointed chain", "number", cp.Number, "checkpoint", cp.Hash, "local", canon) + return + } + log.Error("SCDO checkpoint: local chain CONFLICTS with signed checkpoint, rewinding to checkpoint block", "number", cp.Number, "checkpoint", cp.Hash, "local", canon) + go func() { + if _, err := bc.SetCanonical(block); err != nil { + log.Error("SCDO checkpoint: failed to switch to checkpointed chain", "err", err) + return + } + s.conflict.Store(false) + bc.SetFinalized(block.Header()) + }() +} + +type scdoCPState int + +const ( + scdoCPUnknown scdoCPState = iota // chain tip below checkpoint height + scdoCPContains // chain contains the checkpoint block + scdoCPConflict // chain has another block at checkpoint height +) + +// scdoChainHas reports whether the chain ending at tip contains cp. Cost is +// O(distance from tip to the canonical chain), normally 0 or 1 lookups. +func (bc *BlockChain) scdoChainHas(tip *types.Header, cp *SCDOCheckpoint) scdoCPState { + if tip == nil || tip.Number.Uint64() < cp.Number { + return scdoCPUnknown + } + h := tip + for h != nil && h.Number.Uint64() > cp.Number { + if bc.GetCanonicalHash(h.Number.Uint64()) == h.Hash() { + if bc.GetCanonicalHash(cp.Number) == cp.Hash { + return scdoCPContains + } + return scdoCPConflict + } + h = bc.GetHeader(h.ParentHash, h.Number.Uint64()-1) + } + if h == nil { + return scdoCPUnknown + } + if h.Hash() == cp.Hash { + return scdoCPContains + } + return scdoCPConflict +} + +// scdoForkChoice returns (decided, reorg). decided=false means "no opinion, use TD rules". +func (bc *BlockChain) scdoForkChoice(current, extern *types.Header) (bool, bool) { + s := bc.scdo() + if s == nil { + return false, false + } + cp := s.latest.Load() + if cp == nil { + return false, false + } + cur := bc.scdoChainHas(current, cp) + var ext scdoCPState + if extern.ParentHash == current.Hash() && extern.Number.Uint64() > cp.Number { + ext = cur // plain extension of the current head + } else { + ext = bc.scdoChainHas(extern, cp) + } + switch { + case cur == scdoCPContains && ext != scdoCPContains: + s.rejected.Add(1) + log.Warn("Reorg disallowed by SCDO signed checkpoint 🔒", "checkpoint.bno", cp.Number, "checkpoint.hash", cp.Hash, + "current.bno", current.Number, "current.hash", current.Hash(), "proposed.bno", extern.Number, "proposed.hash", extern.Hash()) + return true, false + case cur == scdoCPConflict && ext == scdoCPContains: + log.Warn("Switching to chain containing SCDO signed checkpoint", "checkpoint.bno", cp.Number, "proposed.bno", extern.Number) + return true, true + } + return false, false +} + +// scdoCheckBlock rejects a block at a checkpointed height with a different hash. +func (bc *BlockChain) scdoCheckBlock(header *types.Header) error { + s := bc.scdo() + if s == nil { + return nil + } + cp := s.latest.Load() + if cp == nil || header.Number.Uint64() != cp.Number || header.Hash() == cp.Hash { + return nil + } + s.rejected.Add(1) + log.Warn("Block rejected by SCDO signed checkpoint 🔒", "number", cp.Number, "checkpoint", cp.Hash, "block", header.Hash()) + return fmt.Errorf("%w: number %d have %s checkpoint %s", ErrSCDOCheckpointMismatch, cp.Number, header.Hash().Hex(), cp.Hash.Hex()) +} diff --git a/core/scdo_checkpoint_test.go b/core/scdo_checkpoint_test.go new file mode 100644 index 0000000..d70461a --- /dev/null +++ b/core/scdo_checkpoint_test.go @@ -0,0 +1,76 @@ +package core + +import ( + "errors" + "testing" + + "github.com/ethereum/go-ethereum/common" + "github.com/ethereum/go-ethereum/crypto" +) + +func scdoTestCP(t testing.TB) (*SCDOCheckpoint, *SCDOCheckpointPolicy) { + key, _ := crypto.GenerateKey() + cp := &SCDOCheckpoint{Version: 1, ChainID: 5680, Number: 1200, Hash: common.HexToHash("0x1234")} + if err := SignSCDOCheckpoint(cp, func(h []byte) ([]byte, error) { return crypto.Sign(h, key) }); err != nil { + t.Fatal(err) + } + return cp, &SCDOCheckpointPolicy{ChainID: 5680, Signers: []common.Address{crypto.PubkeyToAddress(key.PublicKey)}, Threshold: 1} +} + +func TestSCDOCheckpointVerify(t *testing.T) { + cp, pol := scdoTestCP(t) + if _, err := pol.Verify(cp); err != nil { + t.Fatalf("valid checkpoint rejected: %v", err) + } + // tampered hash / number / chainId + for name, mut := range map[string]func(c *SCDOCheckpoint){ + "hash": func(c *SCDOCheckpoint) { c.Hash = common.HexToHash("0x9999") }, + "number": func(c *SCDOCheckpoint) { c.Number++ }, + "sigbyte": func(c *SCDOCheckpoint) { c.Signatures[0][5] ^= 1 }, + } { + cc := &SCDOCheckpoint{Version: 1, ChainID: cp.ChainID, Number: cp.Number, Hash: cp.Hash} + cc.Signatures = append(cc.Signatures, append([]byte{}, cp.Signatures[0]...)) + mut(cc) + if _, err := pol.Verify(cc); err == nil { + t.Fatalf("%s: tampered checkpoint accepted", name) + } + } + c2 := *cp + c2.ChainID = 1 + if _, err := pol.Verify(&c2); !errors.Is(err, ErrSCDOCheckpointChainID) { + t.Fatalf("chainId: %v", err) + } + // signed by a non-authorised key + other, _ := scdoTestCP(t) + if _, err := pol.Verify(other); !errors.Is(err, ErrSCDOCheckpointSigner) { + t.Fatalf("unknown signer: %v", err) + } + // threshold 2 with one signature + pol2 := *pol + pol2.Signers = append(pol2.Signers, common.HexToAddress("0x01")) + pol2.Threshold = 2 + if _, err := pol2.Verify(cp); !errors.Is(err, ErrSCDOCheckpointThreshold) { + t.Fatalf("threshold: %v", err) + } + // duplicate signature does not count twice + c3 := *cp + c3.Signatures = append(c3.Signatures, cp.Signatures[0]) + if _, err := pol2.Verify(&c3); !errors.Is(err, ErrSCDOCheckpointThreshold) { + t.Fatalf("duplicate sig counted twice: %v", err) + } + // no signature + c4 := *cp + c4.Signatures = nil + if _, err := pol.Verify(&c4); !errors.Is(err, ErrSCDOCheckpointNoSig) { + t.Fatalf("nosig: %v", err) + } +} + +func BenchmarkSCDOCheckpointVerify(b *testing.B) { + cp, pol := scdoTestCP(b) + for i := 0; i < b.N; i++ { + if _, err := pol.Verify(cp); err != nil { + b.Fatal(err) + } + } +} diff --git a/eth/backend.go b/eth/backend.go index de71f94..8d154db 100644 --- a/eth/backend.go +++ b/eth/backend.go @@ -76,6 +76,7 @@ type Ethereum struct { txPool *txpool.TxPool blockchain *core.BlockChain + scdoCPSvc *scdoCheckpointService // SCDO signed checkpoint distribution (nil = disabled) handler *handler ethDialCandidates enode.Iterator snapDialCandidates enode.Iterator @@ -260,6 +261,10 @@ func New(stack *node.Node, config *ethconfig.Config) (*Ethereum, error) { } } + if err := eth.setupSCDOCheckpoints(config.SCDOCheckpointSigners, config.SCDOCheckpointThreshold, config.SCDOCheckpointURLs, config.SCDOCheckpointInterval); err != nil { + return nil, err + } + if config.BlobPool.Datadir != "" { config.BlobPool.Datadir = stack.ResolvePath(config.BlobPool.Datadir) } @@ -382,6 +387,9 @@ func (s *Ethereum) APIs() []rpc.API { }, { Namespace: "net", Service: s.netRPCService, + }, { + Namespace: "scdo", + Service: &SCDOCheckpointAPI{s}, }, }...) } @@ -588,6 +596,9 @@ func (s *Ethereum) Start() error { } // Start the networking layer and the light server if requested s.handler.Start(maxPeers) + if s.scdoCPSvc != nil { + s.scdoCPSvc.start() + } return nil } @@ -598,6 +609,9 @@ func (s *Ethereum) Stop() error { s.ethDialCandidates.Close() s.snapDialCandidates.Close() s.handler.Stop() + if s.scdoCPSvc != nil { + s.scdoCPSvc.stop() + } // Then stop everything else. s.bloomIndexer.Close() diff --git a/eth/ethconfig/config.go b/eth/ethconfig/config.go index 881feee..b89a587 100644 --- a/eth/ethconfig/config.go +++ b/eth/ethconfig/config.go @@ -224,6 +224,12 @@ type Config struct { // When this value is *true, ECBP100 will not (ever) be disabled; when *false, it will never be enabled. ECBP1100NoDisable *bool `toml:",omitempty"` + // SCDO signed checkpoints (empty signer list = disabled). + SCDOCheckpointSigners []common.Address `toml:",omitempty"` + SCDOCheckpointThreshold int `toml:",omitempty"` + SCDOCheckpointURLs []string `toml:",omitempty"` + SCDOCheckpointInterval time.Duration `toml:",omitempty"` + // OverrideShanghai (TODO: remove after the fork) OverrideShanghai *uint64 `toml:",omitempty"` diff --git a/eth/handler.go b/eth/handler.go index de87048..f0785c7 100644 --- a/eth/handler.go +++ b/eth/handler.go @@ -99,6 +99,7 @@ type handlerConfig struct { } type handler struct { + scdoBans *scdoPeerBans // SCDO: temporary bans for checkpoint-conflicting peers networkID uint64 forkFilter forkid.Filter // Fork ID filter, constant across the lifetime of the node @@ -143,6 +144,7 @@ func newHandler(config *handlerConfig) (*handler, error) { config.EventMux = new(event.TypeMux) // Nicety initialization for tests } h := &handler{ + scdoBans: newSCDOPeerBans(), networkID: config.Network, forkFilter: forkid.NewFilter(config.Chain), eventMux: config.EventMux, @@ -355,6 +357,11 @@ func (h *handler) runEthPeer(peer *eth.Peer, handler eth.Handler) error { return p2p.DiscQuitting } defer h.decHandlers() + // SCDO: refuse peers that recently served checkpoint-conflicting chains. + if h.scdoBans.banned(peer.ID(), scdoIPOf(peer.RemoteAddr())) { + peer.Log().Debug("Rejecting peer banned for SCDO checkpoint conflicts") + return p2p.DiscUselessPeer + } // If the peer has a `snap` extension, wait for it to connect so we can have // a uniform initialization/teardown mechanism diff --git a/eth/peerset.go b/eth/peerset.go index 4cd8bf6..e8959c4 100644 --- a/eth/peerset.go +++ b/eth/peerset.go @@ -258,6 +258,26 @@ func (ps *peerSet) peerWithHighestTD() *eth.Peer { return bestPeer } +// peerWithHighestTDFiltered is peerWithHighestTD skipping peers for which skip returns true. +func (ps *peerSet) peerWithHighestTDFiltered(skip func(*eth.Peer) bool) *eth.Peer { + ps.lock.RLock() + defer ps.lock.RUnlock() + + var ( + bestPeer *eth.Peer + bestTd *big.Int + ) + for _, p := range ps.peers { + if skip != nil && skip(p.Peer) { + continue + } + if _, td, _ := p.Head(); bestPeer == nil || td.Cmp(bestTd) > 0 { + bestPeer, bestTd = p.Peer, td + } + } + return bestPeer +} + // close disconnects all peers. func (ps *peerSet) close() { ps.lock.Lock() diff --git a/eth/scdo_checkpoint.go b/eth/scdo_checkpoint.go new file mode 100644 index 0000000..4e8a308 --- /dev/null +++ b/eth/scdo_checkpoint.go @@ -0,0 +1,259 @@ +// Copyright 2026 9Y9 PTY LTD (SCDO shard0). LGPL-3.0, same as go-ethereum. +// +// Distribution of SCDO signed checkpoints: poll one or more HTTPS/JSON (or +// file://) sources, and accept pushes over RPC (scdo_submitCheckpoint). Every +// checkpoint is verified against the configured signer list, so a source or +// relay never needs to be trusted; any node can mirror the latest checkpoint. + +package eth + +import ( + "encoding/json" + "errors" + "fmt" + "io" + "net/http" + "net/url" + "os" + "strings" + "sync" + "time" + + "github.com/ethereum/go-ethereum/common" + "github.com/ethereum/go-ethereum/core" + "github.com/ethereum/go-ethereum/log" +) + +const scdoCheckpointMaxBody = 64 * 1024 + +type scdoCheckpointService struct { + eth *Ethereum + urls []string + interval time.Duration + client *http.Client + quit chan struct{} + wg sync.WaitGroup + + mu sync.Mutex + lastFetch time.Time + lastError string + accepted uint64 + invalid uint64 +} + +func newSCDOCheckpointService(eth *Ethereum, urls []string, interval time.Duration) *scdoCheckpointService { + if interval <= 0 { + interval = 30 * time.Second + } + return &scdoCheckpointService{ + eth: eth, urls: urls, interval: interval, + client: &http.Client{Timeout: 10 * time.Second}, + quit: make(chan struct{}), + } +} + +func (s *scdoCheckpointService) start() { + s.wg.Add(1) + go s.loop() +} + +func (s *scdoCheckpointService) stop() { + close(s.quit) + s.wg.Wait() +} + +func (s *scdoCheckpointService) loop() { + defer s.wg.Done() + heads := make(chan core.ChainHeadEvent, 16) + sub := s.eth.blockchain.SubscribeChainHeadEvent(heads) + defer sub.Unsubscribe() + s.pollAll() + t := time.NewTicker(s.interval) + defer t.Stop() + for { + select { + case <-t.C: + s.pollAll() + case <-heads: + s.eth.blockchain.SCDOCheckpointMaintain() + case <-sub.Err(): + return + case <-s.quit: + return + } + } +} + +func (s *scdoCheckpointService) pollAll() { + for _, u := range s.urls { + cp, err := s.fetch(u) + if err == nil { + _, err = s.submit(cp, u) + } + s.mu.Lock() + s.lastFetch = time.Now() + if err != nil { + s.lastError = fmt.Sprintf("%s: %v", u, err) + } else { + s.lastError = "" + } + s.mu.Unlock() + if err != nil && !errors.Is(err, core.ErrSCDOCheckpointOld) { + log.Debug("SCDO checkpoint fetch", "url", u, "err", err) + } + } +} + +func (s *scdoCheckpointService) fetch(src string) (*core.SCDOCheckpoint, error) { + var body []byte + u, err := url.Parse(src) + if err != nil { + return nil, err + } + switch u.Scheme { + case "file": + f, err := os.Open(u.Path) + if err != nil { + return nil, err + } + defer f.Close() + body, err = io.ReadAll(io.LimitReader(f, scdoCheckpointMaxBody)) + if err != nil { + return nil, err + } + case "http", "https": + resp, err := s.client.Get(src) + if err != nil { + return nil, err + } + defer resp.Body.Close() + if resp.StatusCode != http.StatusOK { + return nil, fmt.Errorf("http status %d", resp.StatusCode) + } + body, err = io.ReadAll(io.LimitReader(resp.Body, scdoCheckpointMaxBody)) + if err != nil { + return nil, err + } + default: + return nil, fmt.Errorf("unsupported checkpoint URL scheme %q", u.Scheme) + } + var cp core.SCDOCheckpoint + if err := json.Unmarshal(body, &cp); err != nil { + return nil, err + } + return &cp, nil +} + +func (s *scdoCheckpointService) submit(cp *core.SCDOCheckpoint, src string) (bool, error) { + ok, err := s.eth.blockchain.AddSCDOCheckpoint(cp) + s.mu.Lock() + defer s.mu.Unlock() + switch { + case err == nil && ok: + s.accepted++ + case err != nil && !errors.Is(err, core.ErrSCDOCheckpointOld): + s.invalid++ + log.Warn("SCDO checkpoint rejected", "source", src, "number", cp.Number, "hash", cp.Hash, "err", err) + } + return ok, err +} + +// SCDOCheckpointAPI is the "scdo" RPC namespace. +type SCDOCheckpointAPI struct{ e *Ethereum } + +// GetCheckpoint returns the latest accepted signed checkpoint (null if none). +func (api *SCDOCheckpointAPI) GetCheckpoint() *core.SCDOCheckpoint { + return api.e.blockchain.SCDOCheckpoint() +} + +// SubmitResult is returned by scdo_submitCheckpoint. +type SubmitResult struct { + Accepted bool `json:"accepted"` + Error string `json:"error,omitempty"` +} + +// SubmitCheckpoint verifies and adopts a checkpoint. Safe to expose: only +// checkpoints signed by the configured signers are accepted. +func (api *SCDOCheckpointAPI) SubmitCheckpoint(cp core.SCDOCheckpoint) SubmitResult { + var ( + ok bool + err error + ) + if api.e.scdoCPSvc != nil { + ok, err = api.e.scdoCPSvc.submit(&cp, "rpc") + } else { + ok, err = api.e.blockchain.AddSCDOCheckpoint(&cp) + } + if err != nil { + return SubmitResult{Accepted: false, Error: err.Error()} + } + return SubmitResult{Accepted: ok} +} + +// Status reports the checkpoint subsystem state. +func (api *SCDOCheckpointAPI) Status() map[string]interface{} { + bc := api.e.blockchain + out := map[string]interface{}{"enabled": false} + pol := bc.SCDOCheckpointPolicy() + if pol == nil { + return out + } + rejected, conflict := bc.SCDOCheckpointStats() + out["enabled"] = true + out["chainId"] = pol.ChainID + out["signers"] = pol.Signers + out["threshold"] = pol.Threshold + out["rejectedBlocksOrReorgs"] = rejected + out["conflict"] = conflict + cp := bc.SCDOCheckpoint() + out["checkpoint"] = cp + if cp != nil { + local := bc.GetCanonicalHash(cp.Number) + out["localHashAtCheckpoint"] = local + out["onCheckpointedChain"] = local == cp.Hash + } + if fin := bc.CurrentFinalBlock(); fin != nil { + out["finalized"] = fin.Number.Uint64() + } + out["head"] = bc.CurrentBlock().Number.Uint64() + if api.e.handler != nil && api.e.handler.scdoBans != nil { + out["peerBans"] = api.e.handler.scdoBans.list() + } + if s := api.e.scdoCPSvc; s != nil { + s.mu.Lock() + out["sources"] = s.urls + out["pollInterval"] = s.interval.String() + out["lastFetch"] = s.lastFetch.Unix() + out["lastError"] = s.lastError + out["acceptedCount"] = s.accepted + out["invalidCount"] = s.invalid + s.mu.Unlock() + } + return out +} + +// setupSCDOCheckpoints is called from New() when signers are configured. +func (s *Ethereum) setupSCDOCheckpoints(signers []common.Address, threshold int, urls []string, interval time.Duration) error { + if len(signers) == 0 { + return nil + } + if threshold == 0 { + threshold = 1 + } + chainID := s.blockchain.Config().GetChainID() + if chainID == nil { + return errors.New("SCDO checkpoints need a chainId") + } + pol := &core.SCDOCheckpointPolicy{ChainID: chainID.Uint64(), Signers: signers, Threshold: threshold} + if err := s.blockchain.EnableSCDOCheckpoints(pol); err != nil { + return err + } + var clean []string + for _, u := range urls { + if u = strings.TrimSpace(u); u != "" { + clean = append(clean, u) + } + } + s.scdoCPSvc = newSCDOCheckpointService(s, clean, interval) + return nil +} diff --git a/eth/scdo_peerban.go b/eth/scdo_peerban.go new file mode 100644 index 0000000..b4d545e --- /dev/null +++ b/eth/scdo_peerban.go @@ -0,0 +1,145 @@ +// Copyright 2026 9Y9 PTY LTD (SCDO shard0). LGPL-3.0, same as go-ethereum. +// +// Temporary bans for peers that serve chains conflicting with the SCDO signed +// checkpoint. Without this, a peer that keeps reconnecting with a heavier but +// checkpoint-conflicting chain triggers sync over and over, and a sync attempt +// against an already-dropped peer waits for the downloader timeout (~60 s), +// pausing local mining. Bans are by node ID from the first offence and by IP +// from the second offence from the same IP (loopback is never IP-banned). +// Durations escalate 10 min, 20 min, 40 min ... capped at 24 h. + +package eth + +import ( + "net" + "sort" + "strings" + "sync" + "time" + + "github.com/ethereum/go-ethereum/core" +) + +const ( + scdoBanBase = 10 * time.Minute + scdoBanMax = 24 * time.Hour + scdoBanForget = 48 * time.Hour // offence counters are forgotten after this much quiet time +) + +type scdoBan struct { + count int + until time.Time + last time.Time +} + +type scdoPeerBans struct { + mu sync.Mutex + byID map[string]*scdoBan + byIP map[string]*scdoBan + now func() time.Time +} + +func newSCDOPeerBans() *scdoPeerBans { + return &scdoPeerBans{byID: map[string]*scdoBan{}, byIP: map[string]*scdoBan{}, now: time.Now} +} + +func scdoBanDuration(count int) time.Duration { + d := scdoBanBase + for i := 1; i < count && d < scdoBanMax; i++ { + d *= 2 + } + if d > scdoBanMax { + d = scdoBanMax + } + return d +} + +func scdoIPOf(addr net.Addr) string { + if addr == nil { + return "" + } + host, _, err := net.SplitHostPort(addr.String()) + if err != nil { + host = addr.String() + } + return host +} + +func scdoBanIPAllowed(ip string) bool { + p := net.ParseIP(ip) + return p != nil && !p.IsLoopback() && !p.IsUnspecified() +} + +// isSCDOCheckpointConflict reports whether a sync/import error was caused by a +// chain conflicting with the signed checkpoint (the downloader flattens error chains). +func isSCDOCheckpointConflict(err error) bool { + return err != nil && strings.Contains(err.Error(), core.ErrSCDOCheckpointMismatch.Error()) +} + +func (b *scdoPeerBans) bump(m map[string]*scdoBan, key string, now time.Time) *scdoBan { + e := m[key] + if e == nil || now.Sub(e.last) > scdoBanForget { + e = &scdoBan{} + m[key] = e + } + e.count++ + e.last = now + return e +} + +// note records an offence and returns the ID ban duration. +func (b *scdoPeerBans) note(id, ip string) time.Duration { + b.mu.Lock() + defer b.mu.Unlock() + now := b.now() + e := b.bump(b.byID, id, now) + d := scdoBanDuration(e.count) + e.until = now.Add(d) + if scdoBanIPAllowed(ip) { + ipe := b.bump(b.byIP, ip, now) + if ipe.count >= 2 { + ipe.until = now.Add(scdoBanDuration(ipe.count - 1)) + } + } + return d +} + +// banned reports whether the peer (by ID or IP) is currently banned. +func (b *scdoPeerBans) banned(id, ip string) bool { + b.mu.Lock() + defer b.mu.Unlock() + now := b.now() + if e := b.byID[id]; e != nil && now.Before(e.until) { + return true + } + if ip != "" { + if e := b.byIP[ip]; e != nil && now.Before(e.until) { + return true + } + } + return false +} + +// list returns the active bans (for scdo_status). +func (b *scdoPeerBans) list() []map[string]interface{} { + b.mu.Lock() + defer b.mu.Unlock() + now := b.now() + var out []map[string]interface{} + for k, e := range b.byID { + if now.Before(e.until) { + id := k + if len(id) > 16 { + id = id[:16] + } + out = append(out, map[string]interface{}{"id": id, "offences": e.count, "until": e.until.Unix()}) + } + } + for k, e := range b.byIP { + if now.Before(e.until) { + out = append(out, map[string]interface{}{"ip": k, "offences": e.count, "until": e.until.Unix()}) + } + } + sort.Slice(out, func(i, j int) bool { return out[i]["until"].(int64) < out[j]["until"].(int64) }) + return out +} diff --git a/eth/scdo_peerban_test.go b/eth/scdo_peerban_test.go new file mode 100644 index 0000000..f67be6d --- /dev/null +++ b/eth/scdo_peerban_test.go @@ -0,0 +1,72 @@ +package eth + +import ( + "errors" + "fmt" + "testing" + "time" + + "github.com/ethereum/go-ethereum/core" +) + +func TestSCDOPeerBans(t *testing.T) { + now := time.Unix(1_800_000_000, 0) + b := newSCDOPeerBans() + b.now = func() time.Time { return now } + + if b.banned("a", "1.2.3.4") { + t.Fatal("banned before any offence") + } + if d := b.note("a", "1.2.3.4"); d != 10*time.Minute { + t.Fatalf("first ban %v", d) + } + if !b.banned("a", "") { + t.Fatal("id not banned after offence") + } + if b.banned("b", "1.2.3.4") { + t.Fatal("IP banned after only one offence") + } + // second offence from same IP with a new node key -> IP banned too + b.note("b", "1.2.3.4") + if !b.banned("c", "1.2.3.4") { + t.Fatal("IP not banned after second offence") + } + // escalation for the same id + now = now.Add(11 * time.Minute) + if b.banned("a", "9.9.9.9") { + t.Fatal("ban did not expire") + } + if d := b.note("a", "9.9.9.9"); d != 20*time.Minute { + t.Fatalf("second ban %v", d) + } + for i := 0; i < 20; i++ { + b.note("a", "") + } + if d := scdoBanDuration(30); d != 24*time.Hour { + t.Fatalf("cap %v", d) + } + // counters forgotten after quiet period + now = now.Add(72 * time.Hour) + if d := b.note("a", ""); d != 10*time.Minute { + t.Fatalf("not forgotten: %v", d) + } + // loopback never IP-banned + b.note("x", "127.0.0.1") + b.note("y", "127.0.0.1") + if b.banned("z", "127.0.0.1") { + t.Fatal("loopback IP banned") + } + if len(b.list()) == 0 { + t.Fatal("list empty") + } +} + +func TestSCDOCheckpointConflictDetect(t *testing.T) { + wrapped := fmt.Errorf("retrieved hash chain is invalid: %v", fmt.Errorf("%w: number 80", core.ErrSCDOCheckpointMismatch)) + if !isSCDOCheckpointConflict(wrapped) { + t.Fatal("flattened downloader error not detected") + } + if isSCDOCheckpointConflict(errors.New("timeout")) || isSCDOCheckpointConflict(nil) { + t.Fatal("false positive") + } +} diff --git a/eth/sync.go b/eth/sync.go index 8a103e5..7883a7d 100644 --- a/eth/sync.go +++ b/eth/sync.go @@ -214,7 +214,9 @@ func (cs *chainSyncer) nextSyncOp() *chainSyncOp { // We have enough peers, pick the one with the highest TD, but avoid going // over the terminal total difficulty. Above that we expect the consensus // clients to direct the chain head to sync to. - peer := cs.handler.peers.peerWithHighestTD() + peer := cs.handler.peers.peerWithHighestTDFiltered(func(p *eth.Peer) bool { + return cs.handler.scdoBans.banned(p.ID(), scdoIPOf(p.RemoteAddr())) + }) if peer == nil { return nil } @@ -291,6 +293,11 @@ func (h *handler) doSync(op *chainSyncOp) error { // Run the sync cycle, and disable snap sync if we're past the pivot block err := h.downloader.LegacySync(op.peer.ID(), op.head, op.td, h.chain.Config().GetEthashTerminalTotalDifficulty(), op.mode) if err != nil { + if isSCDOCheckpointConflict(err) { + d := h.scdoBans.note(op.peer.ID(), scdoIPOf(op.peer.RemoteAddr())) + log.Warn("Banning peer for SCDO checkpoint-conflicting chain", "peer", op.peer.ID()[:16], "addr", op.peer.RemoteAddr(), "duration", d) + h.removePeer(op.peer.ID()) + } return err } h.enableSyncedFeatures()