Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion .github/workflows/ci.yml
Original file line number Diff line number Diff line change
Expand Up @@ -17,7 +17,7 @@ jobs:

strategy:
matrix:
go-version: [1.24.x,1.25.x]
go-version: [1.25.x,1.26.x]
os: [ubuntu-latest, macos-latest]

runs-on: ${{ matrix.os }}
Expand Down
5 changes: 3 additions & 2 deletions Makefile
Original file line number Diff line number Diff line change
Expand Up @@ -56,13 +56,14 @@ install:
# Run baseline tests
test: check-testchain-running go-test

# Packages share one testchain and deployment sender, so serialize package binaries.
# Go test short-hand, and skip testing go-ethereum
go-test: test-clean
GOGC=off go test $(TEST_FLAGS) $(MOD_VENDOR) -race -run=$(TEST) `go list ./... | grep -v go-ethereum`
GOGC=off go test $(TEST_FLAGS) $(MOD_VENDOR) -p=1 -race -run=$(TEST) `go list ./... | grep -v go-ethereum`

# Go test short-hand, including go-ethereum
go-test-all: test-clean
GOGC=off go test $(TEST_FLAGS) $(MOD_VENDOR) -run=$(TEST) ./...
GOGC=off go test $(TEST_FLAGS) $(MOD_VENDOR) -p=1 -run=$(TEST) ./...

test-clean:
GOGC=off go clean -testcache
Expand Down
6 changes: 4 additions & 2 deletions ethmonitor/bootstrap.go
Original file line number Diff line number Diff line change
Expand Up @@ -38,8 +38,9 @@ func (c *Chain) bootstrapBlocks(blocks Blocks) error {
return nil
}

if len(blocks) == 1 {
if len(blocks) == 1 && blocks[0].Event != Added {

Copy link
Copy Markdown
Member Author

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 ?

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

A single Added block needs to go through c.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 for Added would leave the bootstrap head untracked and its copies unable to observe a later removal.

The Event != Added condition 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.

TestBlockCanonicalStateBootstrap covers one and three blocks, both directly and through JSON. Those tests passed under -race in this check.

c.blocks = blocks.Copy()
c.blocks[0].canonicalState = nil
return nil
}

Expand All @@ -51,7 +52,7 @@ func (c *Chain) bootstrapBlocks(blocks Blocks) error {

for _, b := range blocks {
if b.Event == Added {
err := c.push(b)
_, err := c.push(b)
if err != nil {
return fmt.Errorf("ethmonitor: bootstrap failed to build canonical chain: %w", err)
}
Expand Down Expand Up @@ -101,5 +102,6 @@ func (b *Block) UnmarshalJSON(data []byte) error {
b.Event = s.Event
b.Logs = s.Logs
b.OK = s.OK
b.canonicalState = nil
return nil
}
131 changes: 131 additions & 0 deletions ethmonitor/cache_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,131 @@
package ethmonitor_test

import (
"context"
"encoding/json"
"errors"
"fmt"
"math/big"
"strings"
"sync/atomic"
"testing"
"time"

"github.com/0xsequence/ethkit/ethmonitor"
"github.com/0xsequence/ethkit/ethmonitor/internal/mocks"
"github.com/0xsequence/ethkit/go-ethereum"
"github.com/0xsequence/ethkit/go-ethereum/common"
"github.com/0xsequence/ethkit/go-ethereum/core/types"
memcache "github.com/goware/cachestore-mem"
cachestore "github.com/goware/cachestore2"
"github.com/stretchr/testify/require"
"go.uber.org/mock/gomock"
)

// Fail before invoking the getter, as an unreachable shared cache would.
type outageBackend struct {
cachestore.Backend
down atomic.Bool
failures atomic.Int64
}

func (b *outageBackend) GetOrSetWithLockEx(ctx context.Context, key string, getter func(context.Context, string) (any, error), ttl time.Duration) (any, error) {
if b.down.Load() {
b.failures.Add(1)
return nil, errors.New("cache unavailable")
}
return b.Backend.GetOrSetWithLockEx(ctx, key, getter, ttl)
}

// Give each fake block a nonzero bloom and an actual log to verify that
// outage recovery preserves logs, rather than just marking blocks ready.
type cacheTestProvider struct {
*fakeProvider
log types.Log
}

func (p *cacheTestProvider) RawBlockByNumber(ctx context.Context, num *big.Int) (json.RawMessage, error) {
payload, err := p.fakeProvider.RawBlockByNumber(ctx, num)
if err != nil {
return nil, err
}
bloom := types.CreateBloom(&types.Receipt{Logs: []*types.Log{&p.log}})
return json.RawMessage(strings.Replace(string(payload), strings.Repeat("0", 512), common.Bytes2Hex(bloom[:]), 1)), nil
}

func (p *cacheTestProvider) RawFilterLogs(ctx context.Context, q ethereum.FilterQuery) (json.RawMessage, error) {
if err := p.chain.wait(ctx); err != nil {
return nil, err
}
block, ok := p.chain.byHash(*q.BlockHash)
if !ok {
return nil, ethereum.NotFound
}
log := p.log
log.BlockHash, log.BlockNumber = block.hash, block.num
return json.Marshal([]types.Log{log})
}

func TestMonitorCacheOutageWithLogs(t *testing.T) {
for _, concurrency := range []int{0, 2} {
t.Run(fmt.Sprintf("prefetch=%d", concurrency), func(t *testing.T) {
chain := newFakeChain(1000, 40, 0)
chain.hold(1020)
provider := &cacheTestProvider{
fakeProvider: &fakeProvider{MockRawInterface: mocks.NewMockRawInterface(gomock.NewController(t)), chain: chain},
log: types.Log{Address: common.HexToAddress("0x1234"), Topics: []common.Hash{common.HexToHash("0xabcd")}, Data: []byte{1, 2, 3}},
}
mem, err := memcache.NewBackend(512)
require.NoError(t, err)
backend := &outageBackend{Backend: mem}
backend.down.Store(true)
opts := ethmonitor.DefaultOptions
opts.PollingInterval = 5 * time.Millisecond
opts.Timeout = time.Second
opts.WithLogs = true
opts.BlockRetentionLimit = 5 // Queue capacity is 10; publish 20 blocks while down.
opts.StartBlockNumber = big.NewInt(1000)
opts.PrefetchConcurrency = concurrency
opts.CacheBackend = backend
monitor, err := ethmonitor.NewMonitor(provider, opts)
require.NoError(t, err)
sub := monitor.Subscribe()
defer sub.Unsubscribe()
done := runMonitorForTest(t, monitor)
timer := time.NewTimer(3 * time.Second)
defer timer.Stop()
next := uint64(1000)
for next <= 1039 {
select {
case blocks := <-sub.Blocks():
for _, block := range blocks {
require.Equal(t, next, block.NumberU64())
require.Equal(t, ethmonitor.Added, block.Event)
require.True(t, block.OK)
require.NotZero(t, block.Bloom())
log := provider.log
log.BlockHash, log.BlockNumber = block.Hash(), next
require.Equal(t, []types.Log{log}, block.Logs)
next++
if next == 1020 {
require.True(t, backend.down.Load())
require.GreaterOrEqual(t, backend.failures.Load(), int64(40))
backend.down.Store(false)
chain.release(1020)
}
}
case err := <-done:
t.Fatalf("Run returned during cache outage/recovery: %v (head=%v)", err, monitor.LatestBlockNum())
case <-timer.C:
t.Fatalf("monitor stopped publishing at block %d", next-1)
}
}
// After recovery, ordinary reads fill the cache again without a restart.
key := ethmonitor.CacheKeyBlockLogs(big.NewInt(1), chain.hashAt(1039), nil)
_, found, err := backend.Get(context.Background(), key)
require.NoError(t, err)
require.True(t, found)
require.True(t, monitor.IsRunning())
})
}
}
179 changes: 179 additions & 0 deletions ethmonitor/canonical_state_test.go
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)
}
}
Loading
Loading