Skip to content

Commit 83fe452

Browse files
committed
fix(evmreader): EVM Reader is ready only when polling sucessfully
1 parent 0ad1d8c commit 83fe452

7 files changed

Lines changed: 110 additions & 54 deletions

File tree

cmd/cartesi-rollups-advancer/root/root.go

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -125,7 +125,7 @@ func run(cmd *cobra.Command, args []string) {
125125
}
126126

127127
supCfg := &service.SupervisorConfigs{
128-
BaseConfigs: service.BaseConfigs{Name: name, Logger: logger},
128+
BaseConfigs: service.BaseConfigs{Name: name, Logger: logger},
129129
EnableSignalHandling: true,
130130
TelemetryCreate: true,
131131
TelemetryAddress: cfg.AdvancerTelemetryAddress,

internal/evmreader/evmreader.go

Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -110,6 +110,9 @@ func (r *Service) Tick(ctx context.Context) (bool, error) {
110110
return false, err
111111
}
112112

113+
now := time.Now()
114+
r.lastSuccessfulPoll.Store(&now)
115+
113116
if blockNumber != r.lastBlockNumber.Load() {
114117
r.lastBlockNumber.Store(blockNumber)
115118
r.Logger.Info("Got new block header", "block", blockNumber, "policy", r.defaultBlock)
@@ -122,6 +125,10 @@ func (r *Service) Tick(ctx context.Context) (bool, error) {
122125
return false, nil
123126
}
124127

128+
func (r *Service) Ready() bool {
129+
return time.Since(*r.lastSuccessfulPoll.Load()) < r.pollingMaxWait
130+
}
131+
125132
func (r *Service) processBlockHead(
126133
ctx context.Context,
127134
blockNumber uint64,

internal/evmreader/evmreader_test.go

Lines changed: 66 additions & 38 deletions
Original file line numberDiff line numberDiff line change
@@ -55,31 +55,26 @@ func (s *EvmReaderSuite) SetupTest() {
5555
s.inputBox = newMockInputBox().SetupDefaultBehavior()
5656
s.contractFactory = newMockAdapterFactory().SetupDefaultBehavior(s.applicationContract1, s.applicationContract2, s.inputBox)
5757

58-
s.evmReader = &Service{
59-
client: s.client,
60-
repository: s.repository,
61-
defaultBlock: DefaultBlock_Latest,
62-
inputReaderEnabled: true,
63-
hasEnabledApps: true,
64-
adapterFactory: s.contractFactory,
65-
}
66-
6758
logLevel, err := config.GetLogLevel()
6859
s.Require().NoError(err)
6960

70-
serviceArgs := &service.TickServiceConfigs{
71-
BaseConfigs: service.BaseConfigs{
72-
Name: "evm-reader",
73-
LogLevel: logLevel,
61+
evmReader, err := Create(s.T().Context(), &CreateInfo{
62+
Config: config.EvmreaderConfig{
63+
LogLevel: logLevel,
64+
BlockchainHttpRetryMaxWait: 200 * time.Millisecond,
65+
BlockchainDefaultBlock: DefaultBlock_Latest,
66+
FeatureInputReaderEnabled: true,
67+
EvmReaderPollingInterval: 100 * time.Millisecond,
7468
},
75-
PollInterval: 100 * time.Millisecond,
76-
}
77-
err = service.InitTickServiceTemplate(&s.evmReader.TickServiceTemplate, serviceArgs, s.evmReader)
69+
Repository: s.repository,
70+
EthClient: s.client,
71+
AdapterFactory: s.contractFactory,
72+
})
7873
s.Require().NoError(err)
79-
80-
s.evmReader.resolver = newApplicationAdapterResolver(s.evmReader.Logger, s.contractFactory)
74+
s.evmReader = evmReader.(*Service)
8175

8276
supCfg := service.SupervisorConfigs{
77+
BaseConfigs: service.BaseConfigs{LogLevel: logLevel},
8378
Factories: []service.FactoryFunction{
8479
func(context.Context, *service.Supervisor) (service.SupervisedService, error) {
8580
return s.evmReader, nil
@@ -132,68 +127,101 @@ func wasntNotified(ch <-chan struct{}) bool {
132127
}
133128

134129
// Service tests
135-
func (s *EvmReaderSuite) TestItStopsWhenSupervisorIsStoppedAfterFirstHeader() {
130+
func (s *EvmReaderSuite) TestItStopsWhenContextIsCanceled() {
136131
called := newCallNotification(s.client.EnqueueNewHead(100))
137132

133+
ctx, cancel := context.WithCancel(s.T().Context())
134+
138135
done := make(chan struct{})
139136
go func() {
140-
err := s.supervisor.Serve(s.T().Context())
141-
s.Require().NoError(err)
137+
err := s.evmReader.Serve(ctx)
138+
s.Require().ErrorIs(err, context.Canceled)
142139
close(done)
143140
}()
144141

145142
s.Require().True(waitNotification(called), "evmreader did not read new header")
146143

147-
s.supervisor.Stop(false)
144+
cancel()
148145

149146
s.Require().True(waitNotification(done), "evmreader did not stop after context cancelation")
150147
}
151148

152149
func (s *EvmReaderSuite) TestReadyReflectsServeLifecycle() {
153150
called := newCallNotification(s.client.EnqueueNewHead(100))
154151

155-
s.Require().False(s.supervisor.Ready())
152+
s.Require().False(s.evmReader.Ready())
153+
154+
ctx, cancel := context.WithCancel(s.T().Context())
156155

157156
done := make(chan struct{})
158157
go func() {
159-
err := s.supervisor.Serve(s.T().Context())
160-
s.Require().NoError(err)
158+
err := s.evmReader.Serve(ctx)
159+
s.Require().ErrorIs(err, context.Canceled)
161160
close(done)
162161
}()
163162

164163
s.Require().True(waitNotification(called))
165-
s.Require().True(s.supervisor.Ready())
166-
s.Require().True(wasntNotified(done))
164+
s.Require().True(s.evmReader.Ready())
167165

168-
s.supervisor.Stop(false)
166+
s.Require().True(wasntNotified(done))
167+
cancel()
169168
s.Require().True(waitNotification(done))
170-
s.Require().False(s.supervisor.Ready())
169+
170+
s.Require().True(s.evmReader.Ready())
171+
time.Sleep(s.evmReader.pollingMaxWait)
172+
s.Require().False(s.evmReader.Ready())
171173
}
172174

173-
func (s *EvmReaderSuite) TestReadyDoesNotDependOnPollingSuccess() {
175+
func (s *EvmReaderSuite) TestNotReadyWhilePollingFails() {
174176
var hdr *types.Header
175177
called := newCallNotification(s.client.On("HeaderByNumber",
176178
mock.Anything,
177179
mock.Anything,
180+
).Return(hdr, errors.New("transient connection error")))
181+
182+
s.Require().False(s.evmReader.Ready())
183+
184+
go func() {
185+
err := s.evmReader.Serve(s.T().Context())
186+
s.Require().ErrorIs(err, context.Canceled)
187+
}()
188+
189+
s.Require().True(waitNotification(called))
190+
s.Require().False(s.evmReader.Ready())
191+
}
192+
193+
func (s *EvmReaderSuite) TestReadyIsRecoveredOnPollingSuccess() {
194+
var hdr *types.Header
195+
failCalled := newCallNotification(s.client.On("HeaderByNumber",
196+
mock.Anything,
197+
mock.Anything,
178198
).Return(hdr, errors.New("transient connection error")).Once())
199+
called := newCallNotification(s.client.EnqueueNewHead(100))
179200

180-
s.Require().False(s.supervisor.Ready())
201+
s.Require().False(s.evmReader.Ready())
202+
203+
ctx, cancel := context.WithCancel(s.T().Context())
181204

182205
done := make(chan struct{})
183206
go func() {
184-
err := s.supervisor.Serve(s.T().Context())
185-
s.Require().NoError(err)
207+
err := s.evmReader.Serve(ctx)
208+
s.Require().ErrorIs(err, context.Canceled)
186209
close(done)
187210
}()
188211

189-
s.Require().True(waitNotification(called))
190-
s.Require().True(s.supervisor.Ready())
191-
s.Require().True(wasntNotified(done))
212+
s.Require().True(waitNotification(failCalled))
213+
s.Require().False(s.evmReader.Ready())
192214

193-
s.supervisor.Stop(false)
215+
s.Require().True(waitNotification(called))
216+
s.Require().True(s.evmReader.Ready())
194217

218+
s.Require().True(wasntNotified(done))
219+
cancel()
195220
s.Require().True(waitNotification(done))
196-
s.Require().False(s.supervisor.Ready())
221+
222+
s.Require().True(s.evmReader.Ready())
223+
time.Sleep(s.evmReader.pollingMaxWait)
224+
s.Require().False(s.evmReader.Ready())
197225
}
198226

199227
func (s *EvmReaderSuite) TestTickScansWithServiceContext() {

internal/evmreader/mocks_test.go

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -76,6 +76,7 @@ func newMockEthClient() *MockEthClient {
7676
}
7777

7878
func (m *MockEthClient) SetupDefaultBehavior() *MockEthClient {
79+
m.On("ChainID", mock.Anything).Return(big.NewInt(0), nil)
7980
return m
8081
}
8182

@@ -189,6 +190,10 @@ func newMockRepository() *MockRepository {
189190
}
190191

191192
func (m *MockRepository) SetupDefaultBehavior() *MockRepository {
193+
m.On("LoadNodeConfigRaw", mock.Anything, EvmReaderConfigKey).
194+
Return(([]byte)(nil), time.Time{}, time.Time{}, repository.ErrNotFound)
195+
m.On("SaveNodeConfigRaw", mock.Anything, EvmReaderConfigKey, mock.Anything).
196+
Return(nil)
192197

193198
apps := copyApplications(applications)
194199
m.On("ListApplications",

internal/evmreader/service.go

Lines changed: 27 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -10,6 +10,7 @@ import (
1010
"log/slog"
1111
"math/big"
1212
"sync/atomic"
13+
"time"
1314

1415
"github.com/cartesi/rollups-node/internal/config"
1516
. "github.com/cartesi/rollups-node/internal/model"
@@ -20,10 +21,11 @@ import (
2021
)
2122

2223
type CreateInfo struct {
23-
Config config.EvmreaderConfig
24-
Logger *slog.Logger
25-
EthClient *ethclient.Client
26-
Repository EvmReaderRepository
24+
Config config.EvmreaderConfig
25+
Logger *slog.Logger
26+
Repository EvmReaderRepository
27+
EthClient EthClientInterface
28+
AdapterFactory AdapterFactory
2729
}
2830

2931
type Service struct {
@@ -38,6 +40,8 @@ type Service struct {
3840
hasEnabledApps bool
3941
inputReaderEnabled bool
4042
lastBlockNumber atomic.Uint64
43+
lastSuccessfulPoll atomic.Pointer[time.Time]
44+
pollingMaxWait time.Duration
4145
}
4246

4347
const EvmReaderConfigKey = "evm-reader"
@@ -116,15 +120,27 @@ func Create(ctx context.Context, c *CreateInfo) (service.SupervisedService, erro
116120
s.defaultBlock = nodeConfig.DefaultBlock
117121
s.inputReaderEnabled = nodeConfig.InputReaderEnabled
118122
s.hasEnabledApps = true
119-
s.adapterFactory = &DefaultAdapterFactory{
120-
Client: ethClient,
121-
Filter: ethutil.Filter{
122-
MinChunkSize: ethutil.DefaultMinChunkSize,
123-
MaxChunkSize: new(big.Int).SetUint64(c.Config.BlockchainMaxBlockRange),
124-
Logger: s.Logger,
125-
},
123+
124+
if c.AdapterFactory != nil {
125+
s.adapterFactory = c.AdapterFactory
126+
} else {
127+
fullEthClient, ok := ethClient.(*ethclient.Client)
128+
if !ok {
129+
return nil, fmt.Errorf("EthClient must be *ethclient.Client when AdapterFactory is not provided")
130+
}
131+
s.adapterFactory = &DefaultAdapterFactory{
132+
Client: fullEthClient,
133+
Filter: ethutil.Filter{
134+
MinChunkSize: ethutil.DefaultMinChunkSize,
135+
MaxChunkSize: new(big.Int).SetUint64(c.Config.BlockchainMaxBlockRange),
136+
Logger: s.Logger,
137+
},
138+
}
126139
}
140+
127141
s.resolver = newApplicationAdapterResolver(s.Logger, s.adapterFactory)
142+
s.pollingMaxWait = c.Config.BlockchainHttpRetryMaxWait
143+
s.lastSuccessfulPoll.Store(&time.Time{})
128144

129145
s.Logger.Info("Created", "config", c.Config)
130146

internal/prt/service.go

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -21,9 +21,9 @@ import (
2121
)
2222

2323
type CreateInfo struct {
24-
Config config.PrtConfig
25-
Logger *slog.Logger
26-
Repository repository.Repository
24+
Config config.PrtConfig
25+
Logger *slog.Logger
26+
Repository repository.Repository
2727
}
2828

2929
type Service struct {

pkg/service/supervisor_test.go

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -367,7 +367,7 @@ func (s *SupervisorSuite) TestFailedInitialization() {
367367
errorChild := newErrorService("error-child")
368368

369369
supervisor := NewSupervisor(&SupervisorConfigs{
370-
BaseConfigs: BaseConfigs{Name: s.T().Name()},
370+
BaseConfigs: BaseConfigs{Name: s.T().Name()},
371371
Factories: []FactoryFunction{
372372
func(context.Context, *Supervisor) (SupervisedService, error) {
373373
return healthyChild, nil

0 commit comments

Comments
 (0)