diff --git a/storage/pebble/registers.go b/storage/pebble/registers.go index b707d677814..3c490947a15 100644 --- a/storage/pebble/registers.go +++ b/storage/pebble/registers.go @@ -1,6 +1,7 @@ package pebble import ( + "bytes" "encoding/binary" "fmt" "math" @@ -11,6 +12,7 @@ import ( "go.uber.org/atomic" "github.com/onflow/flow-go/model/flow" + "github.com/onflow/flow-go/module/irrecoverable" "github.com/onflow/flow-go/storage" ) @@ -118,8 +120,9 @@ func (s *Registers) Store( // Upon restart, it may be in a state where registers are indexed in pebble for the latest height // but the remaining execution data in badger is not, so we skip the indexing step without throwing an error if height == latestHeight { - // already updated - return nil + // already updated, but verify the entries match what was stored, + // so that divergent data is not silently dropped + return s.verifyStoredEntries(entries, height) } nextHeight := latestHeight + 1 @@ -152,6 +155,30 @@ func (s *Registers) Store( return nil } +// verifyStoredEntries checks that each of the given entries matches the value already +// stored at the given height. It is used when Store is called again for the latest height, +// which can happen when an execution node restarts and re-indexes a height whose registers +// were already indexed. A mismatch means the previously stored value would be silently +// kept while the caller assumes the new value was stored, so it must be an error. +func (s *Registers) verifyStoredEntries(entries flow.RegisterEntries, height uint64) error { + for _, entry := range entries { + stored, err := s.Get(entry.Key, height) + if errors.Is(err, storage.ErrNotFound) { + // a missing entry means divergence from what was previously stored. + // Do not wrap storage.ErrNotFound: Store's caller must not mistake + // this for a benign "not found" from Get. + return fmt.Errorf("register %v was not stored at height %d", entry.Key, height) + } + if err != nil { + return irrecoverable.NewExceptionf("cannot verify stored register %v at height %d: %w", entry.Key, height, err) + } + if !bytes.Equal(stored, entry.Value) { + return fmt.Errorf("register %v at height %d was already stored with a different value", entry.Key, height) + } + } + return nil +} + // LatestHeight Gets the latest height of complete registers available func (s *Registers) LatestHeight() uint64 { return s.latestHeight.Load() diff --git a/storage/pebble/registers_test.go b/storage/pebble/registers_test.go index 7736c27a99d..4ce23862a43 100644 --- a/storage/pebble/registers_test.go +++ b/storage/pebble/registers_test.go @@ -86,6 +86,27 @@ func TestRegisters_Store(t *testing.T) { err = r.Store(entries, height2) require.NoError(t, err) + // same height with a conflicting value must fail instead of being silently dropped + conflicting := flow.RegisterEntries{ + {Key: key1, Value: []byte("different")}, + } + err = r.Store(conflicting, height2) + require.Error(t, err) + + // same height with a register that was not stored must fail, + // and the failure must not be detectable as storage.ErrNotFound + notStored := flow.RegisterEntries{ + {Key: flow.RegisterID{Owner: "owner", Key: "key2"}, Value: []byte("value2")}, + } + err = r.Store(notStored, height2) + require.Error(t, err) + require.NotErrorIs(t, err, storage.ErrNotFound) + + // the originally stored value is unchanged + value1, err := r.Get(key1, height2) + require.NoError(t, err) + require.Equal(t, expectedValue1, value1) + // out of range height4 := uint64(4) err = r.Store(entries, height4)