Skip to content

Commit e5537b8

Browse files
committed
fix: selectively evict header or data from cache on block validation failure
Classify block validation errors as FaultHeader or FaultData so the syncer can drop only the offending side of the (header, data) pair from the DA cache. Previously a single mismatch removed both, discarding a potentially legitimate counterpart.
1 parent 2b53162 commit e5537b8

3 files changed

Lines changed: 250 additions & 21 deletions

File tree

‎block/internal/syncing/syncer.go‎

Lines changed: 47 additions & 21 deletions
Original file line numberDiff line numberDiff line change
@@ -746,13 +746,23 @@ func (s *Syncer) trySyncNextBlockWithState(ctx context.Context, event *common.DA
746746
// here only the previous block needs to be applied to proceed to the verification.
747747
// The header validation must be done before applying the block to avoid executing gibberish
748748
if err := s.ValidateBlock(ctx, currentState, data, header); err != nil {
749-
// remove header as da included from cache
750-
s.cache.RemoveHeaderDAIncluded(headerHash)
751-
s.cache.RemoveDataDAIncluded(data.DACommitment().String())
752-
753-
if errors.Is(err, types.ErrUnexpectedProposer) {
754-
return errors.Join(errInvalidBlock, err)
749+
var vErr *BlockValidationError
750+
switch {
751+
case errors.As(err, &vErr):
752+
switch vErr.Fault {
753+
case FaultHeader:
754+
s.cache.RemoveHeaderDAIncluded(headerHash)
755+
case FaultData:
756+
s.cache.RemoveDataDAIncluded(data.DACommitment().String())
757+
}
758+
case errors.Is(err, errInvalidState):
759+
// State divergence does not point at a specific side of the pair;
760+
// the cached entries stay so an honest counterpart can still pair up.
761+
default:
762+
s.cache.RemoveHeaderDAIncluded(headerHash)
763+
s.cache.RemoveDataDAIncluded(data.DACommitment().String())
755764
}
765+
756766
if !errors.Is(err, errInvalidState) && !errors.Is(err, errInvalidBlock) {
757767
return errors.Join(errInvalidBlock, err)
758768
}
@@ -891,38 +901,54 @@ func (s *Syncer) executeTxsWithRetry(ctx context.Context, rawTxs [][]byte, heade
891901
return coreexecutor.ExecuteResult{}, nil
892902
}
893903

894-
// ValidateBlock validates a synced block
895-
// NOTE: if the header was gibberish and somehow passed all validation prior but the data was correct
896-
// or if the data was gibberish and somehow passed all validation prior but the header was correct
897-
// we are still losing both in the pending event. This should never happen.
904+
// ValidateBlock validates a synced block. It runs header-only checks first
905+
// (signature, proposer, sequence) and only then the pair checks between the
906+
// header and the attached data. Failures are wrapped in BlockValidationError so
907+
// callers can drop the right side of the pair from caches without discarding a
908+
// potentially legitimate counterpart.
898909
func (s *Syncer) ValidateBlock(_ context.Context, currState types.State, data *types.Data, header *types.SignedHeader) error {
899910
// Set custom verifier for aggregator node signature
900911
header.SetCustomVerifierForSyncNode(s.options.SyncNodeSignatureBytesProvider)
901912

902913
if err := header.ValidateBasicWithData(data); err != nil { //nolint:contextcheck // validation API does not accept context
903-
return fmt.Errorf("invalid header: %w", err)
914+
return classifyValidationError(fmt.Errorf("invalid header: %w", err))
904915
}
905916

906917
if err := currState.AssertExpectedProposer(header); err != nil {
907-
return errors.Join(errInvalidBlock, err)
918+
return errors.Join(errInvalidBlock, &BlockValidationError{Fault: FaultHeader, Err: err})
908919
}
909920

910921
if err := currState.AssertValidForNextState(header, data); err != nil {
911-
if isExternalBlockValidationError(err) {
912-
return errors.Join(errInvalidBlock, err)
922+
if vErr := classifyValidationError(err); vErr != nil {
923+
return errors.Join(errInvalidBlock, vErr)
913924
}
914925
return errors.Join(errInvalidState, err)
915926
}
916927
return nil
917928
}
918929

919-
func isExternalBlockValidationError(err error) bool {
920-
return errors.Is(err, types.ErrUnexpectedProposer) ||
921-
errors.Is(err, types.ErrInvalidChainID) ||
922-
errors.Is(err, types.ErrInvalidBlockHeight) ||
923-
errors.Is(err, types.ErrInvalidBlockTime) ||
924-
errors.Is(err, types.ErrHeaderDataMismatch) ||
925-
errors.Is(err, types.ErrDataHashMismatch)
930+
// classifyValidationError tags a known external validation error with the side
931+
// at fault. It returns nil for errors that signal state divergence (the caller
932+
// must then classify them as errInvalidState).
933+
func classifyValidationError(err error) *BlockValidationError {
934+
switch {
935+
case errors.Is(err, types.ErrHeaderDataMismatch),
936+
errors.Is(err, types.ErrDataHashMismatch):
937+
return &BlockValidationError{Fault: FaultData, Err: err}
938+
case errors.Is(err, types.ErrUnexpectedProposer),
939+
errors.Is(err, types.ErrInvalidChainID),
940+
errors.Is(err, types.ErrInvalidBlockHeight),
941+
errors.Is(err, types.ErrInvalidBlockTime),
942+
errors.Is(err, types.ErrSignerPubKeyMissing),
943+
errors.Is(err, types.ErrSignerAddressMismatch),
944+
errors.Is(err, types.ErrSignatureEmpty),
945+
errors.Is(err, types.ErrSignatureVerificationFailed),
946+
errors.Is(err, types.ErrProposerAddressMismatch),
947+
errors.Is(err, types.ErrNoProposerAddress):
948+
return &BlockValidationError{Fault: FaultHeader, Err: err}
949+
default:
950+
return nil
951+
}
926952
}
927953

928954
var errMaliciousProposer = errors.New("malicious proposer detected")

‎block/internal/syncing/syncer_test.go‎

Lines changed: 160 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -352,6 +352,166 @@ func TestSyncer_ValidateBlock_RejectsSignerAddressNotDerivedFromPubKey(t *testin
352352
require.Contains(t, err.Error(), "signer address")
353353
}
354354

355+
func TestSyncer_ValidateBlock_ClassifiesFault(t *testing.T) {
356+
expectedAddr, expectedPub, expectedSigner := buildSyncTestSigner(t)
357+
wrongAddr, wrongPub, wrongSigner := buildSyncTestSigner(t)
358+
359+
now := time.Now()
360+
baseState := types.State{
361+
ChainID: "tchain",
362+
InitialHeight: 1,
363+
LastBlockHeight: 1,
364+
LastBlockTime: now,
365+
LastHeaderHash: []byte("last-header-hash"),
366+
AppHash: []byte("app0"),
367+
NextProposerAddress: expectedAddr,
368+
}
369+
370+
makeHeader := func(tb testing.TB, chainID string, proposer []byte, pub crypto.PubKey, signer signerpkg.Signer, appHash []byte, data *types.Data) *types.SignedHeader {
371+
tb.Helper()
372+
_, header := makeSignedHeaderBytes(tb, chainID, 2, proposer, pub, signer, appHash, data, baseState.LastHeaderHash)
373+
return header
374+
}
375+
376+
tests := map[string]struct {
377+
setup func(testing.TB) (*types.SignedHeader, *types.Data)
378+
wantFault ValidationFault
379+
}{
380+
"wrong proposer -> header fault": {
381+
setup: func(tb testing.TB) (*types.SignedHeader, *types.Data) {
382+
data := makeData(baseState.ChainID, 2, 1)
383+
data.Metadata.Time = uint64(now.Add(time.Second).UnixNano())
384+
return makeHeader(tb, baseState.ChainID, wrongAddr, wrongPub, wrongSigner, baseState.AppHash, data), data
385+
},
386+
wantFault: FaultHeader,
387+
},
388+
"wrong chain id -> header fault": {
389+
setup: func(tb testing.TB) (*types.SignedHeader, *types.Data) {
390+
data := makeData("other-chain", 2, 1)
391+
data.Metadata.Time = uint64(now.Add(time.Second).UnixNano())
392+
return makeHeader(tb, "other-chain", expectedAddr, expectedPub, expectedSigner, baseState.AppHash, data), data
393+
},
394+
wantFault: FaultHeader,
395+
},
396+
"data hash mismatch -> data fault": {
397+
setup: func(tb testing.TB) (*types.SignedHeader, *types.Data) {
398+
headerData := makeData(baseState.ChainID, 2, 1)
399+
headerData.Metadata.Time = uint64(now.Add(time.Second).UnixNano())
400+
header := makeHeader(tb, baseState.ChainID, expectedAddr, expectedPub, expectedSigner, baseState.AppHash, headerData)
401+
attached := makeData(baseState.ChainID, 2, 2)
402+
attached.Metadata.Time = headerData.Metadata.Time
403+
return header, attached
404+
},
405+
wantFault: FaultData,
406+
},
407+
"header data metadata mismatch -> data fault": {
408+
setup: func(tb testing.TB) (*types.SignedHeader, *types.Data) {
409+
headerData := makeData(baseState.ChainID, 2, 1)
410+
headerData.Metadata.Time = uint64(now.Add(time.Second).UnixNano())
411+
header := makeHeader(tb, baseState.ChainID, expectedAddr, expectedPub, expectedSigner, baseState.AppHash, headerData)
412+
attached := &types.Data{
413+
Metadata: &types.Metadata{
414+
ChainID: baseState.ChainID,
415+
Height: 3,
416+
Time: headerData.Metadata.Time,
417+
},
418+
Txs: headerData.Txs,
419+
}
420+
return header, attached
421+
},
422+
wantFault: FaultData,
423+
},
424+
}
425+
426+
for name, tc := range tests {
427+
t.Run(name, func(t *testing.T) {
428+
header, data := tc.setup(t)
429+
s := &Syncer{logger: zerolog.Nop(), options: common.DefaultBlockOptions()}
430+
err := s.ValidateBlock(t.Context(), baseState, data, header)
431+
require.Error(t, err)
432+
433+
var vErr *BlockValidationError
434+
require.ErrorAs(t, err, &vErr, "expected *BlockValidationError")
435+
assert.Equal(t, tc.wantFault, vErr.Fault, "fault classification")
436+
})
437+
}
438+
}
439+
440+
func TestSyncer_TrySyncNextBlock_SelectiveCacheCleanup(t *testing.T) {
441+
expectedAddr, expectedPub, expectedSigner := buildSyncTestSigner(t)
442+
wrongAddr, wrongPub, wrongSigner := buildSyncTestSigner(t)
443+
444+
now := time.Now()
445+
baseState := types.State{
446+
ChainID: "tchain",
447+
InitialHeight: 1,
448+
LastBlockHeight: 1,
449+
LastBlockTime: now,
450+
LastHeaderHash: []byte("last-header-hash"),
451+
AppHash: []byte("app0"),
452+
NextProposerAddress: expectedAddr,
453+
}
454+
455+
makeSyncer := func(tb testing.TB) (*Syncer, cache.Manager) {
456+
tb.Helper()
457+
ds := dssync.MutexWrap(datastore.NewMapDatastore())
458+
st := store.New(ds)
459+
cm, err := cache.NewManager(config.DefaultConfig(), st, zerolog.Nop())
460+
require.NoError(tb, err)
461+
return &Syncer{cache: cm, logger: zerolog.Nop(), options: common.DefaultBlockOptions()}, cm
462+
}
463+
464+
tests := map[string]struct {
465+
event func(testing.TB) common.DAHeightEvent
466+
wantHeaderInCache bool
467+
wantDataInCache bool
468+
}{
469+
"header fault keeps data in cache": {
470+
event: func(tb testing.TB) common.DAHeightEvent {
471+
data := makeData(baseState.ChainID, 2, 1)
472+
data.Metadata.Time = uint64(now.Add(time.Second).UnixNano())
473+
_, header := makeSignedHeaderBytes(tb, baseState.ChainID, 2, wrongAddr, wrongPub, wrongSigner, baseState.AppHash, data, baseState.LastHeaderHash)
474+
return common.DAHeightEvent{Header: header, Data: data, Source: common.SourceDA}
475+
},
476+
wantHeaderInCache: false,
477+
wantDataInCache: true,
478+
},
479+
"data fault keeps header in cache": {
480+
event: func(tb testing.TB) common.DAHeightEvent {
481+
headerData := makeData(baseState.ChainID, 2, 1)
482+
headerData.Metadata.Time = uint64(now.Add(time.Second).UnixNano())
483+
_, header := makeSignedHeaderBytes(tb, baseState.ChainID, 2, expectedAddr, expectedPub, expectedSigner, baseState.AppHash, headerData, baseState.LastHeaderHash)
484+
attached := makeData(baseState.ChainID, 2, 2)
485+
attached.Metadata.Time = headerData.Metadata.Time
486+
return common.DAHeightEvent{Header: header, Data: attached, Source: common.SourceDA}
487+
},
488+
wantHeaderInCache: true,
489+
wantDataInCache: false,
490+
},
491+
}
492+
493+
for name, tc := range tests {
494+
t.Run(name, func(t *testing.T) {
495+
s, cm := makeSyncer(t)
496+
event := tc.event(t)
497+
headerHash := event.Header.Hash().String()
498+
dataHash := event.Data.DACommitment().String()
499+
500+
// Seed the cache so we can observe what gets removed.
501+
cm.SetHeaderDAIncluded(headerHash, 1, event.Header.Height())
502+
cm.SetDataDAIncluded(dataHash, 1, event.Header.Height())
503+
504+
err := s.trySyncNextBlockWithState(t.Context(), &event, baseState)
505+
require.Error(t, err)
506+
507+
_, headerStillIncluded := cm.GetHeaderDAIncludedByHash(headerHash)
508+
_, dataStillIncluded := cm.GetDataDAIncludedByHash(dataHash)
509+
assert.Equal(t, tc.wantHeaderInCache, headerStillIncluded, "header cache presence")
510+
assert.Equal(t, tc.wantDataInCache, dataStillIncluded, "data cache presence")
511+
})
512+
}
513+
}
514+
355515
func TestSyncer_ApplyBlockPersistsExecutionNextProposer(t *testing.T) {
356516
addr, _, _ := buildSyncTestSigner(t)
357517
execNext := []byte("execution-next-proposer")
Lines changed: 43 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,43 @@
1+
package syncing
2+
3+
import "fmt"
4+
5+
// ValidationFault identifies which side of a (header, data) pair is responsible
6+
// for a block validation failure. Callers use it to decide what to evict from
7+
// caches so a legitimate counterpart is not discarded together with the bad one.
8+
type ValidationFault int
9+
10+
const (
11+
// FaultHeader marks the header as invalid on its own (signature, proposer,
12+
// chain id, height, sequence). The data attached to the event may still be
13+
// legitimate and pair with a different valid header.
14+
FaultHeader ValidationFault = iota
15+
16+
// FaultData marks the data as the suspect. It is used when the header
17+
// passed all header-only checks but the data does not match the header
18+
// (mismatched DataHash or metadata).
19+
FaultData
20+
)
21+
22+
// BlockValidationError wraps the underlying validation error with a fault tag.
23+
type BlockValidationError struct {
24+
Fault ValidationFault
25+
Err error
26+
}
27+
28+
func (e *BlockValidationError) Error() string {
29+
return fmt.Sprintf("block validation failed (%s): %v", e.Fault, e.Err)
30+
}
31+
32+
func (e *BlockValidationError) Unwrap() error { return e.Err }
33+
34+
func (f ValidationFault) String() string {
35+
switch f {
36+
case FaultHeader:
37+
return "header"
38+
case FaultData:
39+
return "data"
40+
default:
41+
return "unknown"
42+
}
43+
}

0 commit comments

Comments
 (0)