Repository navigation
ethmonitor: parallel block prefetch for fast chains #219
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Changes from 6 commits
cf8c4c1
1c35532
987262d
247c2e1
77bbef8
95199ae
a43e30f
62a096a
0dba4f8
8233c60
1fc3e78
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,179 @@ | ||
| package ethmonitor | ||
|
|
||
| import ( | ||
| "encoding/json" | ||
| "fmt" | ||
| "math/big" | ||
| "testing" | ||
|
|
||
| "github.com/0xsequence/ethkit/go-ethereum/common" | ||
| "github.com/0xsequence/ethkit/go-ethereum/core/types" | ||
| "github.com/stretchr/testify/require" | ||
| ) | ||
|
|
||
| func canonicalTestBlock(num int64) *Block { | ||
| return &Block{ | ||
| Block: types.NewBlockWithHeader(&types.Header{ | ||
| Number: big.NewInt(num), | ||
| BlockHash: common.BigToHash(big.NewInt(num)), | ||
| ParentHash: common.BigToHash(big.NewInt(num - 1)), | ||
| Time: uint64(num), | ||
| }), | ||
| Event: Added, | ||
| OK: true, | ||
| } | ||
| } | ||
|
|
||
| func pushCanonicalTestBlock(t *testing.T, chain *Chain, input *Block) *Block { | ||
| t.Helper() | ||
| chain.push(input) | ||
| block := chain.Head() | ||
| require.NotNil(t, block) | ||
| require.Equal(t, input.Hash(), block.Hash()) | ||
| return block | ||
| } | ||
|
|
||
| // The interface keeps the before-fix proof executable without the new method. | ||
| func canonicalTestState(t *testing.T, block *Block) (uint64, bool) { | ||
| t.Helper() | ||
| state, ok := any(block).(interface{ CanonicalState() (uint64, bool) }) | ||
| if !ok { | ||
| return 0, false | ||
| } | ||
| return state.CanonicalState() | ||
| } | ||
|
|
||
| func TestBlockCanonicalStateRetention(t *testing.T) { | ||
| for _, depth := range []int64{1, 3} { | ||
| t.Run(fmt.Sprintf("evictionDepth=%d", depth), func(t *testing.T) { | ||
| chain := newChain(10, false) | ||
| block := pushCanonicalTestBlock(t, chain, canonicalTestBlock(100)) | ||
| copy := chain.Blocks().Copy()[0] | ||
| incarnation, canonical := canonicalTestState(t, block) | ||
| require.Positive(t, incarnation) | ||
| require.True(t, canonical) | ||
| for num := int64(101); num <= 109+depth; num++ { | ||
| pushCanonicalTestBlock(t, chain, canonicalTestBlock(num)) | ||
| } | ||
| require.Nil(t, chain.GetBlock(block.Hash())) | ||
| for _, snapshot := range []*Block{block, copy} { | ||
| got, canonical := canonicalTestState(t, snapshot) | ||
| require.Equal(t, incarnation, got) | ||
| require.True(t, canonical, "retention eviction was mistaken for removal") | ||
| } | ||
| }) | ||
| } | ||
| } | ||
|
|
||
| func TestBlockCanonicalStateReadoption(t *testing.T) { | ||
| chain := newChain(10, false) | ||
| input := canonicalTestBlock(100) | ||
| block := pushCanonicalTestBlock(t, chain, input) | ||
| shallow := *block | ||
| copy := chain.Blocks().Copy()[0] | ||
| incarnation, canonical := canonicalTestState(t, block) | ||
| require.Positive(t, incarnation) | ||
| require.True(t, canonical) | ||
| removed := *chain.pop() | ||
| removed.Event = Removed | ||
| for _, snapshot := range []*Block{block, &shallow, copy, &removed} { | ||
| got, canonical := canonicalTestState(t, snapshot) | ||
| require.Equal(t, incarnation, got) | ||
| require.False(t, canonical, "snapshot did not observe the actual removal") | ||
| } | ||
|
|
||
| // Reusing either the original input or a removed snapshot must create a fresh | ||
| // owned incarnation without mutating old Added/Removed event copies. | ||
| for _, reused := range []*Block{input, copy} { | ||
| fresh := pushCanonicalTestBlock(t, chain, reused) | ||
| got, canonical := canonicalTestState(t, fresh) | ||
| require.Greater(t, got, incarnation) | ||
| require.True(t, canonical) | ||
| for _, snapshot := range []*Block{block, &shallow, copy, &removed} { | ||
| got, canonical := canonicalTestState(t, snapshot) | ||
| require.Equal(t, incarnation, got) | ||
| require.False(t, canonical, "fresh readoption revived an old event") | ||
| } | ||
| chain.pop() | ||
| } | ||
| } | ||
|
|
||
| func TestBlockCanonicalStateBootstrap(t *testing.T) { | ||
| for _, count := range []int{1, 3} { | ||
| for _, serialized := range []bool{false, true} { | ||
| t.Run(fmt.Sprintf("blocks=%d/JSON=%v", count, serialized), func(t *testing.T) { | ||
| inputs := make(Blocks, count) | ||
| for i := range inputs { | ||
| inputs[i] = canonicalTestBlock(int64(100 + i)) | ||
| } | ||
| chain := newChain(10, true) | ||
| if serialized { | ||
| data, err := json.Marshal(inputs) | ||
| require.NoError(t, err) | ||
| require.NoError(t, chain.BootstrapFromBlocksJSON(data)) | ||
| } else { | ||
| require.NoError(t, chain.BootstrapFromBlocks(inputs)) | ||
| } | ||
| for _, block := range chain.Blocks().Copy() { | ||
| incarnation, canonical := canonicalTestState(t, block) | ||
| require.Positive(t, incarnation) | ||
| require.True(t, canonical) | ||
| data, err := json.Marshal(block) | ||
| require.NoError(t, err) | ||
| require.NotContains(t, string(data), "incarnation") | ||
| require.NoError(t, json.Unmarshal(data, block)) | ||
| incarnation, canonical = canonicalTestState(t, block) | ||
| require.Zero(t, incarnation, "serialized state preserved runtime ownership") | ||
| require.False(t, canonical) | ||
| } | ||
| }) | ||
| } | ||
| } | ||
| } | ||
|
|
||
| func TestBlockCanonicalStateConcurrentRemoval(t *testing.T) { | ||
| chain := newChain(10, false) | ||
| input := canonicalTestBlock(100) | ||
| block := pushCanonicalTestBlock(t, chain, input) | ||
| incarnation, _ := canonicalTestState(t, block) | ||
| require.Positive(t, incarnation) | ||
| state := any(block).(interface{ CanonicalState() (uint64, bool) }) | ||
| start, done := make(chan struct{}), make(chan struct{}) | ||
| result := make(chan error, 1) | ||
| go func() { | ||
| close(start) | ||
| removed := false | ||
| for { | ||
| got, canonical := state.CanonicalState() | ||
| if got != incarnation || (removed && canonical) { | ||
| result <- fmt.Errorf("old incarnation changed or revived: id=%d canonical=%v", got, canonical) | ||
| return | ||
| } | ||
| removed = removed || !canonical | ||
| select { | ||
| case <-done: | ||
| result <- nil | ||
| return | ||
| default: | ||
| } | ||
| } | ||
| }() | ||
| <-start | ||
| for i := 0; i < 100; i++ { | ||
| chain.pop() | ||
| pushCanonicalTestBlock(t, chain, input) | ||
| } | ||
| close(done) | ||
| require.NoError(t, <-result) | ||
| got, canonical := state.CanonicalState() | ||
| require.Equal(t, incarnation, got) | ||
| require.False(t, canonical) | ||
| } | ||
|
|
||
| func TestBlockCanonicalStateUntracked(t *testing.T) { | ||
| for _, block := range []*Block{nil, canonicalTestBlock(100)} { | ||
| incarnation, canonical := canonicalTestState(t, block) | ||
| require.Zero(t, incarnation) | ||
| require.False(t, canonical) | ||
| } | ||
| } |
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -7,6 +7,7 @@ import ( | |
| "math/big" | ||
| "slices" | ||
| "sync" | ||
| "sync/atomic" | ||
|
|
||
| "github.com/0xsequence/ethkit/go-ethereum" | ||
| "github.com/0xsequence/ethkit/go-ethereum/common" | ||
|
|
@@ -28,6 +29,7 @@ type Chain struct { | |
|
|
||
| mu sync.RWMutex | ||
| averageBlockTime float64 // in seconds | ||
| lastIncarnation uint64 | ||
| } | ||
|
|
||
| func newChain(retentionLimit int, bootstrapMode bool) *Chain { | ||
|
|
@@ -59,7 +61,7 @@ func newChain(retentionLimit int, bootstrapMode bool) *Chain { | |
| // } | ||
|
|
||
| // Push to the top of the stack | ||
| func (c *Chain) push(nextBlock *Block) error { | ||
| func (c *Chain) push(nextBlock *Block) (*Block, error) { | ||
| c.mu.Lock() | ||
| defer c.mu.Unlock() | ||
|
|
||
|
|
@@ -70,12 +72,12 @@ func (c *Chain) push(nextBlock *Block) error { | |
|
|
||
| // Assert pointing at prev block | ||
| if nextBlock.ParentHash() != headBlock.Hash() { | ||
| return ErrUnexpectedParentHash | ||
| return nil, ErrUnexpectedParentHash | ||
| } | ||
|
|
||
| // Assert block numbers are in sequence | ||
| if nextBlock.NumberU64() != headBlock.NumberU64()+1 { | ||
| return ErrUnexpectedBlockNumber | ||
| return nil, ErrUnexpectedBlockNumber | ||
| } | ||
|
|
||
| // Update average block time | ||
|
|
@@ -86,14 +88,20 @@ func (c *Chain) push(nextBlock *Block) error { | |
| } | ||
| } | ||
|
|
||
| // Each adoption owns its state so reusing an input cannot revive old events. | ||
| c.lastIncarnation++ | ||
| block := *nextBlock | ||
| block.canonicalState = &blockCanonicalState{incarnation: c.lastIncarnation} | ||
|
Member
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. why do we need block.canonicalState ..? will this add more memory overhead then what we already have?
Member
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Yes, there is a small memory and allocation cost. I measured the layout on amd64: The reason is receipt correctness when the receipt listener falls behind monitor retention. After H is removed and readopted, a delayed This supports the receipt hardening, and applies even with prefetch disabled. There is no additional monitor map retaining every historical block: the state becomes collectible once retained blocks and event copies release it.
Member
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. ok |
||
| block.canonicalState.canonical.Store(true) | ||
|
|
||
| // Add to head of stack | ||
| c.blocks = append(c.blocks, nextBlock) | ||
| c.blocks = append(c.blocks, &block) | ||
| if len(c.blocks) > c.retentionLimit { | ||
| c.blocks[0] = nil | ||
| c.blocks = c.blocks[1:] | ||
| } | ||
|
|
||
| return nil | ||
| return &block, nil | ||
| } | ||
|
|
||
| // Pop from the top of the stack | ||
|
|
@@ -107,6 +115,9 @@ func (c *Chain) pop() *Block { | |
|
|
||
| n := len(c.blocks) - 1 | ||
| block := c.blocks[n] | ||
| if block.canonicalState != nil { | ||
| block.canonicalState.canonical.Store(false) | ||
| } | ||
| c.blocks[n] = nil | ||
| c.blocks = c.blocks[:n] | ||
| return block | ||
|
|
@@ -215,6 +226,9 @@ const ( | |
| Removed | ||
| ) | ||
|
|
||
| // Block contains a monitored block and its event data. | ||
| // Construct values with keyed composite literals: private canonical state makes | ||
| // positional literals unsupported. | ||
| type Block struct { | ||
| *types.Block | ||
|
|
||
|
|
@@ -228,6 +242,26 @@ type Block struct { | |
|
|
||
| // OK flag which represents the block is ready for broadcasting | ||
| OK bool | ||
|
|
||
| canonicalState *blockCanonicalState | ||
|
Member
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. what is
Member
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. It is runtime metadata shared by the monitor's copies of one block adoption:
An
One explicit compatibility tradeoff: adding the private field means external positional |
||
| } | ||
|
|
||
| type blockCanonicalState struct { | ||
| incarnation uint64 | ||
| canonical atomic.Bool | ||
| } | ||
|
|
||
| // CanonicalState reports the monitor-assigned incarnation and whether it has | ||
| // remained canonical without a known removal. Zero means the block is untracked. | ||
| // Retention eviction preserves this state; in-memory copies share removal updates. | ||
| // Incarnations are local to a monitor chain. Serialized blocks are untracked | ||
| // until accepted by a monitor. | ||
| func (b *Block) CanonicalState() (incarnation uint64, canonical bool) { | ||
| if b == nil || b.canonicalState == nil { | ||
| return 0, false | ||
| } | ||
| state := b.canonicalState | ||
| return state.incarnation, state.canonical.Load() | ||
| } | ||
|
|
||
| type Blocks []*Block | ||
|
|
@@ -560,10 +594,11 @@ func (blocks Blocks) Copy() Blocks { | |
| } | ||
|
|
||
| nb[i] = &Block{ | ||
| Block: b.Block, | ||
| Event: b.Event, | ||
| Logs: logs, | ||
| OK: b.OK, | ||
| Block: b.Block, | ||
| Event: b.Event, | ||
| Logs: logs, | ||
| OK: b.OK, | ||
| canonicalState: b.canonicalState, | ||
| } | ||
| } | ||
|
|
||
|
|
||
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
why the
&& blocks[0].Event != Added?There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
A single
Addedblock needs to go throughc.push()now, because that is where the monitor assigns a fresh canonical-state incarnation. Previously every one-block bootstrap took the copy-only shortcut; keeping that shortcut forAddedwould leave the bootstrap head untracked and its copies unable to observe a later removal.The
Event != Addedcondition preserves the existing special handling for a single non-Added input while routing canonical additions through the same adoption path as multi-block bootstrap. It does not change the supplied block hash or fetch anything from the provider.TestBlockCanonicalStateBootstrapcovers one and three blocks, both directly and through JSON. Those tests passed under-racein this check.