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
21 changes: 11 additions & 10 deletions ethmonitor/ethmonitor.go
Original file line number Diff line number Diff line change
Expand Up @@ -869,16 +869,9 @@ func (m *Monitor) filterLogs(ctx context.Context, blockHash common.Hash, topics
if err != nil {
return nil, err
}
if blockBloom != (types.Bloom{}) && (len(logsPayload) == 0 || (len(logsPayload) == 2 && logsPayload[0] == '[' && logsPayload[1] == ']')) {
// If we have no logs and the block bloom is set, then we need to return an error
// as the node is incorrectly telling us the block-logs response is '[]' but in fact
// the block log bloom filter tells us we should be expecting logs. We do this to
// ensure we do not incorrectly cache an empty block-logs response as valid.
return nil, fmt.Errorf("ethmonitor: filterLogs detected empty block-logs response but block bloom is set, ignoring node response")
}
// Validate before caching so a malformed response cannot block log
// backfilling until cache expiry.
fetchedLogs, err = m.unmarshalLogs(logsPayload)
fetchedLogs, err = m.unmarshalLogs(logsPayload, blockBloom)
if err != nil {
return nil, err
}
Expand All @@ -901,7 +894,7 @@ func (m *Monitor) filterLogs(ctx context.Context, blockHash common.Hash, topics
if fetchedLogs != nil {
return fetchedLogs, resp, nil
}
logs, err := m.unmarshalLogs(resp)
logs, err := m.unmarshalLogs(resp, blockBloom)
if err != nil {
// Recover entries cached by peers or older monitors that did not
// validate logs before writing them.
Expand Down Expand Up @@ -1424,11 +1417,19 @@ func (m *Monitor) unmarshalBlock(blockPayload []byte) (*types.Block, error) {
return block, nil
}

func (m *Monitor) unmarshalLogs(logsPayload []byte) ([]types.Log, error) {
func (m *Monitor) unmarshalLogs(logsPayload []byte, blockBloom types.Bloom) ([]types.Log, error) {
var logs []types.Log
err := json.Unmarshal(logsPayload, &logs)
if err != nil {
return nil, err
}
// JSON null decodes without error, but a logs response must be an array.
if logs == nil {
return nil, fmt.Errorf("ethmonitor: filterLogs expected a JSON array of block logs")
}
// Apply the bloom guard to decoded responses, including cache hits.
if len(logs) == 0 && blockBloom != (types.Bloom{}) {
return nil, fmt.Errorf("ethmonitor: filterLogs detected empty block-logs response but block bloom is set, ignoring node response")
}
return logs, nil
}
17 changes: 12 additions & 5 deletions ethmonitor/prefetch_internal_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -40,7 +40,7 @@ func TestPrefetchRejectsInvalidBlockBeforeCaching(t *testing.T) {
}

func TestMonitorRecoversInvalidPrefetchedLogs(t *testing.T) {
for _, payload := range []string{`{}`, `[{}]`} {
for _, payload := range []string{`{}`, `[{}]`, `null`, `[]`, `[ ]`} {
for _, cached := range []bool{false, true} {
t.Run(fmt.Sprintf("payload=%s/cached=%v", payload, cached), func(t *testing.T) {
const first, target, last = uint64(1000), uint64(1002), uint64(1004)
Expand Down Expand Up @@ -114,12 +114,12 @@ func TestMonitorRecoversInvalidPrefetchedLogs(t *testing.T) {
require.NoError(t, monitor.cache.SetEx(context.Background(), key, []byte(payload), time.Hour))
} else {
// Complete the speculative fetch before starting the serial loop,
// so the worker deterministically receives the malformed response.
// so the worker deterministically receives the invalid response.
monitor.prefetch.fetch(context.Background(), prefetchJob{num: target})
require.Equal(t, int64(1), originCalls.Load())
_, found, err := monitor.cache.Get(context.Background(), key)
require.NoError(t, err)
require.False(t, found, "failed prefetch must not cache malformed logs")
require.False(t, found, "failed prefetch must not cache invalid logs")
}

sub := monitor.Subscribe("TestMonitorRecoversInvalidPrefetchedLogs")
Expand Down Expand Up @@ -157,7 +157,7 @@ func TestMonitorRecoversInvalidPrefetchedLogs(t *testing.T) {
}
}
if !cached {
require.GreaterOrEqual(t, originCalls.Load(), int64(2), "malformed prefetch must allow an origin retry")
require.GreaterOrEqual(t, originCalls.Load(), int64(2), "invalid prefetch must allow an origin retry")
} else {
require.Positive(t, originCalls.Load(), "invalid cache entry must allow an origin retry")
}
Expand Down Expand Up @@ -187,7 +187,14 @@ func TestPrefetchPanicStopsWorkers(t *testing.T) {
defer close(done)
monitor.prefetch.run(ctx)
}()
defer func() { cancel(); <-done }()
defer func() {
cancel()
select {
case <-done:
case <-time.After(time.Second):
t.Error("prefetch did not stop after cancellation")
}
}()
select {
case <-done:
case <-time.After(time.Second):
Expand Down
Loading
Loading