Sitelet https://github.com/ethereum-optimism/optimism/commit/88b1a95a5d6732b0103d24808202c35aa199d968
Skip to content

Commit 88b1a95

Browse files
authored
feat: concurrent alt-da requests (#11698)
* feat: initial goroutine blob submission implementation test(batcher): add e2e test for concurrent altda requests doc: add explanation comment for FakeDAServer chore: fix if condition in altda sendTransaction path feat: add maxConcurrentDaRequests config flag + semaphore refactor: batcher to use errgroup for da instead of separate semaphore/waitgroup fix: nil pointer bug after using wrong function after rebase fix: defn of maxConcurrentDaRequests=0 fix: TestBatcherConcurrentAltDARequests chore: remove unneeded if statement around time.Sleep refactor: use TryGo instead of Go to make logic local and easier to read chore: clean up some comments in batcher chore: make batcher shutdown cancel pending altda requests by using shutdownCtx instead of killCtx * chore(batcher): make altda wg wait + log only when useAltDa is true * refactor: batcher altda submission code into its own function * test: refactor batcher e2e test to only count batcher txs * chore: log errors from wait functions * chore: refactor and minimize time that e2e batcher system tests can run * chore: lower timeout duration in test * fix(batcher): maxConcurentDARequests was not being initialized
1 parent 144a775 commit 88b1a95

12 files changed

Lines changed: 378 additions & 109 deletions

File tree

‎op-alt-da/cli.go‎

Lines changed: 42 additions & 14 deletions
Original file line numberDiff line numberDiff line change
@@ -3,15 +3,19 @@ package altda
33
import (
44
"fmt"
55
"net/url"
6+
"time"
67

78
"github.com/urfave/cli/v2"
89
)
910

1011
var (
11-
EnabledFlagName = altDAFlags("enabled")
12-
DaServerAddressFlagName = altDAFlags("da-server")
13-
VerifyOnReadFlagName = altDAFlags("verify-on-read")
14-
DaServiceFlag = altDAFlags("da-service")
12+
EnabledFlagName = altDAFlags("enabled")
13+
DaServerAddressFlagName = altDAFlags("da-server")
14+
VerifyOnReadFlagName = altDAFlags("verify-on-read")
15+
DaServiceFlagName = altDAFlags("da-service")
16+
PutTimeoutFlagName = altDAFlags("put-timeout")
17+
GetTimeoutFlagName = altDAFlags("get-timeout")
18+
MaxConcurrentRequestsFlagName = altDAFlags("max-concurrent-da-requests")
1519
)
1620

1721
// altDAFlags returns the flag names for altDA
@@ -46,20 +50,41 @@ func CLIFlags(envPrefix string, category string) []cli.Flag {
4650
Category: category,
4751
},
4852
&cli.BoolFlag{
49-
Name: DaServiceFlag,
53+
Name: DaServiceFlagName,
5054
Usage: "Use DA service type where commitments are generated by Alt-DA server",
5155
Value: false,
5256
EnvVars: altDAEnvs(envPrefix, "DA_SERVICE"),
5357
Category: category,
5458
},
59+
&cli.DurationFlag{
60+
Name: PutTimeoutFlagName,
61+
Usage: "Timeout for put requests. 0 means no timeout.",
62+
Value: time.Duration(0),
63+
EnvVars: altDAEnvs(envPrefix, "PUT_TIMEOUT"),
64+
},
65+
&cli.DurationFlag{
66+
Name: GetTimeoutFlagName,
67+
Usage: "Timeout for get requests. 0 means no timeout.",
68+
Value: time.Duration(0),
69+
EnvVars: altDAEnvs(envPrefix, "GET_TIMEOUT"),
70+
},
71+
&cli.Uint64Flag{
72+
Name: MaxConcurrentRequestsFlagName,
73+
Usage: "Maximum number of concurrent requests to the DA server",
74+
Value: 1,
75+
EnvVars: altDAEnvs(envPrefix, "MAX_CONCURRENT_DA_REQUESTS"),
76+
},
5577
}
5678
}
5779

5880
type CLIConfig struct {
59-
Enabled bool
60-
DAServerURL string
61-
VerifyOnRead bool
62-
GenericDA bool
81+
Enabled bool
82+
DAServerURL string
83+
VerifyOnRead bool
84+
GenericDA bool
85+
PutTimeout time.Duration
86+
GetTimeout time.Duration
87+
MaxConcurrentRequests uint64
6388
}
6489

6590
func (c CLIConfig) Check() error {
@@ -75,14 +100,17 @@ func (c CLIConfig) Check() error {
75100
}
76101

77102
func (c CLIConfig) NewDAClient() *DAClient {
78-
return &DAClient{url: c.DAServerURL, verify: c.VerifyOnRead, precompute: !c.GenericDA}
103+
return &DAClient{url: c.DAServerURL, verify: c.VerifyOnRead, precompute: !c.GenericDA, getTimeout: c.GetTimeout, putTimeout: c.PutTimeout}
79104
}
80105

81106
func ReadCLIConfig(c *cli.Context) CLIConfig {
82107
return CLIConfig{
83-
Enabled: c.Bool(EnabledFlagName),
84-
DAServerURL: c.String(DaServerAddressFlagName),
85-
VerifyOnRead: c.Bool(VerifyOnReadFlagName),
86-
GenericDA: c.Bool(DaServiceFlag),
108+
Enabled: c.Bool(EnabledFlagName),
109+
DAServerURL: c.String(DaServerAddressFlagName),
110+
VerifyOnRead: c.Bool(VerifyOnReadFlagName),
111+
GenericDA: c.Bool(DaServiceFlagName),
112+
PutTimeout: c.Duration(PutTimeoutFlagName),
113+
GetTimeout: c.Duration(GetTimeoutFlagName),
114+
MaxConcurrentRequests: c.Uint64(MaxConcurrentRequestsFlagName),
87115
}
88116
}

‎op-alt-da/daclient.go‎

Lines changed: 14 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -7,6 +7,7 @@ import (
77
"fmt"
88
"io"
99
"net/http"
10+
"time"
1011
)
1112

1213
// ErrNotFound is returned when the server could not find the input.
@@ -23,10 +24,16 @@ type DAClient struct {
2324
verify bool
2425
// whether commitment is precomputable (only applicable to keccak256)
2526
precompute bool
27+
getTimeout time.Duration
28+
putTimeout time.Duration
2629
}
2730

2831
func NewDAClient(url string, verify bool, pc bool) *DAClient {
29-
return &DAClient{url, verify, pc}
32+
return &DAClient{
33+
url: url,
34+
verify: verify,
35+
precompute: pc,
36+
}
3037
}
3138

3239
// GetInput returns the input data for the given encoded commitment bytes.
@@ -35,7 +42,8 @@ func (c *DAClient) GetInput(ctx context.Context, comm CommitmentData) ([]byte, e
3542
if err != nil {
3643
return nil, fmt.Errorf("failed to create HTTP request: %w", err)
3744
}
38-
resp, err := http.DefaultClient.Do(req)
45+
client := &http.Client{Timeout: c.getTimeout}
46+
resp, err := client.Do(req)
3947
if err != nil {
4048
return nil, err
4149
}
@@ -91,7 +99,8 @@ func (c *DAClient) setInputWithCommit(ctx context.Context, comm CommitmentData,
9199
return fmt.Errorf("failed to create HTTP request: %w", err)
92100
}
93101
req.Header.Set("Content-Type", "application/octet-stream")
94-
resp, err := http.DefaultClient.Do(req)
102+
client := &http.Client{Timeout: c.putTimeout}
103+
resp, err := client.Do(req)
95104
if err != nil {
96105
return err
97106
}
@@ -116,7 +125,8 @@ func (c *DAClient) setInput(ctx context.Context, img []byte) (CommitmentData, er
116125
return nil, fmt.Errorf("failed to create HTTP request: %w", err)
117126
}
118127
req.Header.Set("Content-Type", "application/octet-stream")
119-
resp, err := http.DefaultClient.Do(req)
128+
client := &http.Client{Timeout: c.putTimeout}
129+
resp, err := client.Do(req)
120130
if err != nil {
121131
return nil, err
122132
}

‎op-alt-da/daclient_test.go‎

Lines changed: 2 additions & 36 deletions
Original file line numberDiff line numberDiff line change
@@ -2,48 +2,14 @@ package altda
22

33
import (
44
"context"
5-
"fmt"
65
"math/rand"
7-
"sync"
86
"testing"
97

108
"github.com/ethereum-optimism/optimism/op-service/testlog"
11-
"github.com/ethereum/go-ethereum/common"
129
"github.com/ethereum/go-ethereum/log"
1310
"github.com/stretchr/testify/require"
1411
)
1512

16-
type MemStore struct {
17-
db map[string][]byte
18-
lock sync.RWMutex
19-
}
20-
21-
func NewMemStore() *MemStore {
22-
return &MemStore{
23-
db: make(map[string][]byte),
24-
}
25-
}
26-
27-
// Get retrieves the given key if it's present in the key-value store.
28-
func (s *MemStore) Get(ctx context.Context, key []byte) ([]byte, error) {
29-
s.lock.RLock()
30-
defer s.lock.RUnlock()
31-
32-
if entry, ok := s.db[string(key)]; ok {
33-
return common.CopyBytes(entry), nil
34-
}
35-
return nil, ErrNotFound
36-
}
37-
38-
// Put inserts the given value into the key-value store.
39-
func (s *MemStore) Put(ctx context.Context, key []byte, value []byte) error {
40-
s.lock.Lock()
41-
defer s.lock.Unlock()
42-
43-
s.db[string(key)] = common.CopyBytes(value)
44-
return nil
45-
}
46-
4713
func TestDAClientPrecomputed(t *testing.T) {
4814
store := NewMemStore()
4915
logger := testlog.Logger(t, log.LevelDebug)
@@ -56,7 +22,7 @@ func TestDAClientPrecomputed(t *testing.T) {
5622

5723
cfg := CLIConfig{
5824
Enabled: true,
59-
DAServerURL: fmt.Sprintf("http://%s", server.Endpoint()),
25+
DAServerURL: server.HttpEndpoint(),
6026
VerifyOnRead: true,
6127
}
6228
require.NoError(t, cfg.Check())
@@ -113,7 +79,7 @@ func TestDAClientService(t *testing.T) {
11379

11480
cfg := CLIConfig{
11581
Enabled: true,
116-
DAServerURL: fmt.Sprintf("http://%s", server.Endpoint()),
82+
DAServerURL: server.HttpEndpoint(),
11783
VerifyOnRead: false,
11884
GenericDA: false,
11985
}

‎op-alt-da/damock.go‎

Lines changed: 85 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -4,8 +4,12 @@ import (
44
"context"
55
"errors"
66
"io"
7+
"net/http"
8+
"sync"
9+
"time"
710

811
"github.com/ethereum-optimism/optimism/op-service/eth"
12+
"github.com/ethereum/go-ethereum/common"
913
"github.com/ethereum/go-ethereum/ethdb"
1014
"github.com/ethereum/go-ethereum/ethdb/memorydb"
1115
"github.com/ethereum/go-ethereum/log"
@@ -99,3 +103,84 @@ func (d *AltDADisabled) OnFinalizedHeadSignal(f HeadSignalFn) {
99103
func (d *AltDADisabled) AdvanceL1Origin(ctx context.Context, l1 L1Fetcher, blockId eth.BlockID) error {
100104
return ErrNotEnabled
101105
}
106+
107+
// FakeDAServer is a fake DA server for e2e tests.
108+
// It is a small wrapper around DAServer that allows for setting request latencies,
109+
// to mimic a DA service with slow responses (eg. eigenDA with 10 min batching interval).
110+
type FakeDAServer struct {
111+
*DAServer
112+
putRequestLatency time.Duration
113+
getRequestLatency time.Duration
114+
}
115+
116+
func NewFakeDAServer(host string, port int, log log.Logger) *FakeDAServer {
117+
store := NewMemStore()
118+
fakeDAServer := &FakeDAServer{
119+
DAServer: NewDAServer(host, port, store, log, true),
120+
putRequestLatency: 0,
121+
getRequestLatency: 0,
122+
}
123+
return fakeDAServer
124+
}
125+
126+
func (s *FakeDAServer) HandleGet(w http.ResponseWriter, r *http.Request) {
127+
time.Sleep(s.getRequestLatency)
128+
s.DAServer.HandleGet(w, r)
129+
}
130+
131+
func (s *FakeDAServer) HandlePut(w http.ResponseWriter, r *http.Request) {
132+
time.Sleep(s.putRequestLatency)
133+
s.DAServer.HandlePut(w, r)
134+
}
135+
136+
func (s *FakeDAServer) Start() error {
137+
err := s.DAServer.Start()
138+
if err != nil {
139+
return err
140+
}
141+
// Override the HandleGet/Put method registrations
142+
mux := http.NewServeMux()
143+
mux.HandleFunc("/get/", s.HandleGet)
144+
mux.HandleFunc("/put/", s.HandlePut)
145+
s.httpServer.Handler = mux
146+
return nil
147+
}
148+
149+
func (s *FakeDAServer) SetPutRequestLatency(latency time.Duration) {
150+
s.putRequestLatency = latency
151+
}
152+
153+
func (s *FakeDAServer) SetGetRequestLatency(latency time.Duration) {
154+
s.getRequestLatency = latency
155+
}
156+
157+
type MemStore struct {
158+
db map[string][]byte
159+
lock sync.RWMutex
160+
}
161+
162+
func NewMemStore() *MemStore {
163+
return &MemStore{
164+
db: make(map[string][]byte),
165+
}
166+
}
167+
168+
// Get retrieves the given key if it's present in the key-value store.
169+
func (s *MemStore) Get(ctx context.Context, key []byte) ([]byte, error) {
170+
s.lock.RLock()
171+
defer s.lock.RUnlock()
172+
173+
if entry, ok := s.db[string(key)]; ok {
174+
return common.CopyBytes(entry), nil
175+
}
176+
return nil, ErrNotFound
177+
}
178+
179+
// Put inserts the given value into the key-value store.
180+
func (s *MemStore) Put(ctx context.Context, key []byte, value []byte) error {
181+
s.lock.Lock()
182+
defer s.lock.Unlock()
183+
184+
s.db[string(key)] = common.CopyBytes(value)
185+
return nil
186+
}

‎op-alt-da/daserver.go‎

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -187,8 +187,8 @@ func (d *DAServer) HandlePut(w http.ResponseWriter, r *http.Request) {
187187
}
188188
}
189189

190-
func (b *DAServer) Endpoint() string {
191-
return b.listener.Addr().String()
190+
func (b *DAServer) HttpEndpoint() string {
191+
return fmt.Sprintf("http://%s", b.listener.Addr().String())
192192
}
193193

194194
func (b *DAServer) Stop() error {

0 commit comments

Comments
 (0)