Skip to content

Commit e277538

Browse files
committed
refactor(submitter): concurrent submitter
1 parent 49ef5c9 commit e277538

8 files changed

Lines changed: 413 additions & 511 deletions

block/internal/submitting/da_submitter.go

Lines changed: 397 additions & 234 deletions
Large diffs are not rendered by default.

block/internal/submitting/da_submitter_integration_test.go

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -107,6 +107,8 @@ func TestDASubmitter_SubmitHeadersAndData_MarksInclusionAndUpdatesLastSubmitted(
107107
require.NoError(t, err)
108108
require.NoError(t, daSubmitter.SubmitData(context.Background(), dataList, marshalledData, cm, n, gen))
109109

110+
daSubmitter.Flush()
111+
110112
// After submission, inclusion markers should be set
111113
_, ok := cm.GetHeaderDAIncludedByHeight(1)
112114
assert.True(t, ok)
Lines changed: 0 additions & 270 deletions
Original file line numberDiff line numberDiff line change
@@ -1,271 +1 @@
11
package submitting
2-
3-
import (
4-
"context"
5-
"testing"
6-
"time"
7-
8-
"github.com/rs/zerolog"
9-
"github.com/stretchr/testify/assert"
10-
"github.com/stretchr/testify/mock"
11-
12-
"github.com/evstack/ev-node/block/internal/common"
13-
"github.com/evstack/ev-node/pkg/config"
14-
datypes "github.com/evstack/ev-node/pkg/da/types"
15-
"github.com/evstack/ev-node/pkg/genesis"
16-
"github.com/evstack/ev-node/test/mocks"
17-
)
18-
19-
// helper to build a basic submitter with provided DA mock client and config overrides
20-
func newTestSubmitter(t *testing.T, mockClient *mocks.MockClient, override func(*config.Config)) *DASubmitter {
21-
cfg := config.Config{}
22-
// Keep retries small and backoffs minimal
23-
cfg.DA.BlockTime.Duration = 1 * time.Millisecond
24-
cfg.DA.MaxSubmitAttempts = 3
25-
cfg.DA.SubmitOptions = "opts"
26-
cfg.DA.Namespace = "ns"
27-
cfg.DA.DataNamespace = "ns-data"
28-
if override != nil {
29-
override(&cfg)
30-
}
31-
if mockClient == nil {
32-
mockClient = mocks.NewMockClient(t)
33-
}
34-
mockClient.On("GetHeaderNamespace").Return([]byte(cfg.DA.Namespace)).Maybe()
35-
mockClient.On("GetDataNamespace").Return([]byte(cfg.DA.DataNamespace)).Maybe()
36-
mockClient.On("GetForcedInclusionNamespace").Return([]byte(nil)).Maybe()
37-
mockClient.On("HasForcedInclusionNamespace").Return(false).Maybe()
38-
return NewDASubmitter(mockClient, cfg, genesis.Genesis{} /*options=*/, common.BlockOptions{}, common.NopMetrics(), zerolog.Nop(), nil, nil)
39-
}
40-
41-
func TestSubmitToDA_MempoolRetry_IncreasesGasAndSucceeds(t *testing.T) {
42-
t.Parallel()
43-
44-
client := mocks.NewMockClient(t)
45-
46-
nsBz := datypes.NamespaceFromString("ns").Bytes()
47-
opts := []byte("opts")
48-
var usedGas []float64
49-
50-
client.On("Submit", mock.Anything, mock.Anything, mock.AnythingOfType("float64"), nsBz, opts).
51-
Run(func(args mock.Arguments) {
52-
usedGas = append(usedGas, args.Get(2).(float64))
53-
}).
54-
Return(datypes.ResultSubmit{BaseResult: datypes.BaseResult{Code: datypes.StatusNotIncludedInBlock, SubmittedCount: 0}}).
55-
Once()
56-
57-
ids := [][]byte{[]byte("id1"), []byte("id2"), []byte("id3")}
58-
client.On("Submit", mock.Anything, mock.Anything, mock.AnythingOfType("float64"), nsBz, opts).
59-
Run(func(args mock.Arguments) {
60-
usedGas = append(usedGas, args.Get(2).(float64))
61-
}).
62-
Return(datypes.ResultSubmit{BaseResult: datypes.BaseResult{Code: datypes.StatusSuccess, IDs: ids, SubmittedCount: uint64(len(ids))}}).
63-
Once()
64-
65-
s := newTestSubmitter(t, client, nil)
66-
67-
items := []string{"a", "b", "c"}
68-
marshalledItems := make([][]byte, len(items))
69-
for idx, item := range items {
70-
marshalledItems[idx] = []byte(item)
71-
}
72-
73-
ctx := context.Background()
74-
err := submitToDA[string](
75-
s,
76-
ctx,
77-
items,
78-
marshalledItems,
79-
func(_ []string, _ *datypes.ResultSubmit) {},
80-
"item",
81-
nsBz,
82-
opts,
83-
)
84-
assert.NoError(t, err)
85-
86-
// Sentinel value is preserved on retry
87-
assert.Equal(t, []float64{-1, -1}, usedGas)
88-
}
89-
90-
func TestSubmitToDA_UnknownError_RetriesSameGasThenSucceeds(t *testing.T) {
91-
t.Parallel()
92-
93-
client := mocks.NewMockClient(t)
94-
95-
nsBz := datypes.NamespaceFromString("ns").Bytes()
96-
97-
opts := []byte("opts")
98-
var usedGas []float64
99-
100-
// First attempt: unknown failure -> reasonFailure, gas unchanged for next attempt
101-
client.On("Submit", mock.Anything, mock.Anything, mock.AnythingOfType("float64"), nsBz, opts).
102-
Run(func(args mock.Arguments) { usedGas = append(usedGas, args.Get(2).(float64)) }).
103-
Return(datypes.ResultSubmit{BaseResult: datypes.BaseResult{Code: datypes.StatusError, Message: "boom"}}).
104-
Once()
105-
106-
// Second attempt: same gas, success
107-
ids := [][]byte{[]byte("id1")}
108-
client.On("Submit", mock.Anything, mock.Anything, mock.AnythingOfType("float64"), nsBz, opts).
109-
Run(func(args mock.Arguments) { usedGas = append(usedGas, args.Get(2).(float64)) }).
110-
Return(datypes.ResultSubmit{BaseResult: datypes.BaseResult{Code: datypes.StatusSuccess, IDs: ids, SubmittedCount: uint64(len(ids))}}).
111-
Once()
112-
113-
s := newTestSubmitter(t, client, nil)
114-
115-
items := []string{"x"}
116-
marshalledItems := make([][]byte, len(items))
117-
for idx, item := range items {
118-
marshalledItems[idx] = []byte(item)
119-
}
120-
121-
ctx := context.Background()
122-
err := submitToDA[string](
123-
s,
124-
ctx,
125-
items,
126-
marshalledItems,
127-
func(_ []string, _ *datypes.ResultSubmit) {},
128-
"item",
129-
nsBz,
130-
opts,
131-
)
132-
assert.NoError(t, err)
133-
assert.Equal(t, []float64{-1, -1}, usedGas)
134-
}
135-
136-
func TestSubmitToDA_TooBig_HalvesBatch(t *testing.T) {
137-
t.Parallel()
138-
139-
client := mocks.NewMockClient(t)
140-
141-
nsBz := datypes.NamespaceFromString("ns").Bytes()
142-
143-
opts := []byte("opts")
144-
var batchSizes []int
145-
146-
client.On("Submit", mock.Anything, mock.Anything, mock.Anything, nsBz, opts).
147-
Run(func(args mock.Arguments) {
148-
blobs := args.Get(1).([][]byte)
149-
batchSizes = append(batchSizes, len(blobs))
150-
}).
151-
Return(datypes.ResultSubmit{BaseResult: datypes.BaseResult{Code: datypes.StatusTooBig}}).
152-
Once()
153-
154-
ids := [][]byte{[]byte("id1"), []byte("id2")}
155-
client.On("Submit", mock.Anything, mock.Anything, mock.Anything, nsBz, opts).
156-
Run(func(args mock.Arguments) {
157-
blobs := args.Get(1).([][]byte)
158-
batchSizes = append(batchSizes, len(blobs))
159-
}).
160-
Return(datypes.ResultSubmit{BaseResult: datypes.BaseResult{Code: datypes.StatusSuccess, IDs: ids, SubmittedCount: uint64(len(ids))}}).
161-
Once()
162-
163-
s := newTestSubmitter(t, client, nil)
164-
165-
items := []string{"a", "b", "c", "d"}
166-
marshalledItems := make([][]byte, len(items))
167-
for idx, item := range items {
168-
marshalledItems[idx] = []byte(item)
169-
}
170-
171-
ctx := context.Background()
172-
err := submitToDA[string](
173-
s,
174-
ctx,
175-
items,
176-
marshalledItems,
177-
func(_ []string, _ *datypes.ResultSubmit) {},
178-
"item",
179-
nsBz,
180-
opts,
181-
)
182-
assert.NoError(t, err)
183-
assert.Equal(t, []int{4, 2}, batchSizes)
184-
}
185-
186-
func TestSubmitToDA_SentinelNoGas_PreservesGasAcrossRetries(t *testing.T) {
187-
t.Parallel()
188-
189-
client := mocks.NewMockClient(t)
190-
191-
nsBz := datypes.NamespaceFromString("ns").Bytes()
192-
193-
opts := []byte("opts")
194-
var usedGas []float64
195-
196-
client.On("Submit", mock.Anything, mock.Anything, mock.AnythingOfType("float64"), nsBz, opts).
197-
Run(func(args mock.Arguments) { usedGas = append(usedGas, args.Get(2).(float64)) }).
198-
Return(datypes.ResultSubmit{BaseResult: datypes.BaseResult{Code: datypes.StatusAlreadyInMempool}}).
199-
Once()
200-
201-
ids := [][]byte{[]byte("id1")}
202-
client.On("Submit", mock.Anything, mock.Anything, mock.AnythingOfType("float64"), nsBz, opts).
203-
Run(func(args mock.Arguments) { usedGas = append(usedGas, args.Get(2).(float64)) }).
204-
Return(datypes.ResultSubmit{BaseResult: datypes.BaseResult{Code: datypes.StatusSuccess, IDs: ids, SubmittedCount: uint64(len(ids))}}).
205-
Once()
206-
207-
s := newTestSubmitter(t, client, nil)
208-
209-
items := []string{"only"}
210-
marshalledItems := make([][]byte, len(items))
211-
for idx, item := range items {
212-
marshalledItems[idx] = []byte(item)
213-
}
214-
215-
ctx := context.Background()
216-
err := submitToDA[string](
217-
s,
218-
ctx,
219-
items,
220-
marshalledItems,
221-
func(_ []string, _ *datypes.ResultSubmit) {},
222-
"item",
223-
nsBz,
224-
opts,
225-
)
226-
assert.NoError(t, err)
227-
assert.Equal(t, []float64{-1, -1}, usedGas)
228-
}
229-
230-
func TestSubmitToDA_PartialSuccess_AdvancesWindow(t *testing.T) {
231-
t.Parallel()
232-
233-
client := mocks.NewMockClient(t)
234-
235-
nsBz := datypes.NamespaceFromString("ns").Bytes()
236-
237-
opts := []byte("opts")
238-
var totalSubmitted int
239-
240-
firstIDs := [][]byte{[]byte("id1"), []byte("id2")}
241-
client.On("Submit", mock.Anything, mock.Anything, mock.Anything, nsBz, opts).
242-
Return(datypes.ResultSubmit{BaseResult: datypes.BaseResult{Code: datypes.StatusSuccess, IDs: firstIDs, SubmittedCount: uint64(len(firstIDs))}}).
243-
Once()
244-
245-
secondIDs := [][]byte{[]byte("id3")}
246-
client.On("Submit", mock.Anything, mock.Anything, mock.Anything, nsBz, opts).
247-
Return(datypes.ResultSubmit{BaseResult: datypes.BaseResult{Code: datypes.StatusSuccess, IDs: secondIDs, SubmittedCount: uint64(len(secondIDs))}}).
248-
Once()
249-
250-
s := newTestSubmitter(t, client, nil)
251-
252-
items := []string{"a", "b", "c"}
253-
marshalledItems := make([][]byte, len(items))
254-
for idx, item := range items {
255-
marshalledItems[idx] = []byte(item)
256-
}
257-
258-
ctx := context.Background()
259-
err := submitToDA[string](
260-
s,
261-
ctx,
262-
items,
263-
marshalledItems,
264-
func(submitted []string, _ *datypes.ResultSubmit) { totalSubmitted += len(submitted) },
265-
"item",
266-
nsBz,
267-
opts,
268-
)
269-
assert.NoError(t, err)
270-
assert.Equal(t, 3, totalSubmitted)
271-
}

block/internal/submitting/da_submitter_test.go

Lines changed: 4 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -35,18 +35,15 @@ const (
3535
func setupDASubmitterTest(t *testing.T) (*DASubmitter, store.Store, cache.Manager, *mocks.MockClient, genesis.Genesis) {
3636
t.Helper()
3737

38-
// Create store and cache
3938
ds := sync.MutexWrap(datastore.NewMapDatastore())
4039
st := store.New(ds)
4140
cm, err := cache.NewManager(config.DefaultConfig(), st, zerolog.Nop())
4241
require.NoError(t, err)
4342

44-
// Create config
4543
cfg := config.DefaultConfig()
4644
cfg.DA.Namespace = testHeaderNamespace
4745
cfg.DA.DataNamespace = testDataNamespace
4846

49-
// Mock DA client
5047
mockDA := mocks.NewMockClient(t)
5148
headerNamespace := datypes.NamespaceFromString(cfg.DA.Namespace).Bytes()
5249
dataNamespace := datypes.NamespaceFromString(cfg.DA.DataNamespace).Bytes()
@@ -55,15 +52,13 @@ func setupDASubmitterTest(t *testing.T) (*DASubmitter, store.Store, cache.Manage
5552
mockDA.On("GetForcedInclusionNamespace").Return([]byte(nil)).Maybe()
5653
mockDA.On("HasForcedInclusionNamespace").Return(false).Maybe()
5754

58-
// Create genesis
5955
gen := genesis.Genesis{
6056
ChainID: "test-chain",
6157
InitialHeight: 1,
6258
StartTime: time.Now(),
6359
ProposerAddress: []byte("test-proposer"),
6460
}
6561

66-
// Create DA submitter
6762
daSubmitter := NewDASubmitter(
6863
mockDA,
6964
cfg,
@@ -220,6 +215,8 @@ func TestDASubmitter_SubmitHeaders_Success(t *testing.T) {
220215
require.NoError(t, err)
221216
err = submitter.SubmitHeaders(ctx, headers, marshalledHeaders, cm, signer)
222217
require.NoError(t, err)
218+
submitter.Flush()
219+
submitter.Close()
223220

224221
// Verify headers are marked as DA included
225222
_, ok1 := cm.GetHeaderDAIncludedByHeight(1)
@@ -335,6 +332,8 @@ func TestDASubmitter_SubmitData_Success(t *testing.T) {
335332
require.NoError(t, err)
336333
err = submitter.SubmitData(ctx, signedDataList, marshalledData, cm, signer, gen)
337334
require.NoError(t, err)
335+
submitter.Flush()
336+
submitter.Close()
338337

339338
// Verify data is marked as DA included
340339
_, ok := cm.GetDataDAIncludedByHeight(1)

block/internal/submitting/da_submitter_tracing.go

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -95,3 +95,7 @@ func (t *tracedDASubmitter) SubmitData(ctx context.Context, signedDataList []*ty
9595

9696
return nil
9797
}
98+
99+
func (t *tracedDASubmitter) Close() {
100+
t.inner.Close()
101+
}

block/internal/submitting/da_submitter_tracing_test.go

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -37,6 +37,8 @@ func (m *mockDASubmitterAPI) SubmitData(ctx context.Context, signedDataList []*t
3737
return nil
3838
}
3939

40+
func (m *mockDASubmitterAPI) Close() {}
41+
4042
func setupDASubmitterTrace(t *testing.T, inner DASubmitterAPI) (DASubmitterAPI, *tracetest.SpanRecorder) {
4143
t.Helper()
4244
sr := tracetest.NewSpanRecorder()

block/internal/submitting/submitter.go

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -23,10 +23,10 @@ import (
2323
"github.com/evstack/ev-node/types"
2424
)
2525

26-
// DASubmitterAPI defines minimal methods needed by Submitter for DA submissions.
2726
type DASubmitterAPI interface {
2827
SubmitHeaders(ctx context.Context, headers []*types.SignedHeader, marshalledHeaders [][]byte, cache cache.Manager, signer signer.Signer) error
2928
SubmitData(ctx context.Context, signedDataList []*types.SignedData, marshalledData [][]byte, cache cache.Manager, signer signer.Signer, genesis genesis.Genesis) error
29+
Close()
3030
}
3131

3232
// Submitter handles DA submission and inclusion processing for both sync and aggregator nodes
@@ -136,7 +136,6 @@ func (s *Submitter) Start(ctx context.Context) (err error) {
136136
return err
137137
}
138138

139-
// Start DA submission loop if signer is available (aggregator nodes only)
140139
if s.signer != nil {
141140
s.logger.Info().Msg("starting DA submission loop")
142141
s.wg.Go(s.daSubmissionLoop)
@@ -153,6 +152,7 @@ func (s *Submitter) Stop() error {
153152
if s.cancel != nil {
154153
s.cancel()
155154
}
155+
s.daSubmitter.Close()
156156
// Wait for goroutines to finish with a timeout to prevent hanging
157157
done := make(chan struct{})
158158
go func() {

block/internal/submitting/submitter_test.go

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -437,6 +437,8 @@ func (f *fakeDASubmitter) SubmitData(ctx context.Context, _ []*types.SignedData,
437437
return nil
438438
}
439439

440+
func (f *fakeDASubmitter) Close() {}
441+
440442
// fakeSigner implements signer.Signer with deterministic behavior for tests.
441443
type fakeSigner struct{}
442444

0 commit comments

Comments
 (0)