Skip to content

Commit aaacbde

Browse files
committed
Async DA pull
1 parent e2b0520 commit aaacbde

6 files changed

Lines changed: 514 additions & 44 deletions

File tree

block/internal/common/expected_interfaces.go

Lines changed: 0 additions & 21 deletions
Original file line numberDiff line numberDiff line change
@@ -20,24 +20,3 @@ type Broadcaster[H header.Header[H]] interface {
2020
AppendDAHint(ctx context.Context, daHeight uint64, hashes ...types.Hash) error
2121
GetByHeight(ctx context.Context, height uint64) (H, uint64, error)
2222
}
23-
24-
//
25-
//// Decorator to access the payload type without the container
26-
//type Decorator[H header.Header[H]] struct {
27-
// nested Broadcaster[*sync.DAHeightHintContainer[H]]
28-
//}
29-
//
30-
//func NewDecorator[H header.Header[H]](nested Broadcaster[*sync.DAHeightHintContainer[H]]) Decorator[H] {
31-
// return Decorator[H]{nested: nested}
32-
//}
33-
//
34-
//func (d Decorator[H]) WriteToStoreAndBroadcast(ctx context.Context, payload H, opts ...pubsub.PubOpt) error {
35-
// return d.nested.WriteToStoreAndBroadcast(ctx, &sync.DAHeightHintContainer[H]{Entry: payload}, opts...)
36-
//}
37-
//
38-
//func (d Decorator[H]) AppendDAHint(ctx context.Context, daHeight uint64, hashes ...types.Hash) error {
39-
// return d.nested.AppendDAHint(ctx, daHeight, hashes...)
40-
//}
41-
//func (d Decorator[H]) GetByHeight(ctx context.Context, height uint64) (H, error) {
42-
// panic("not implemented")
43-
//}
Lines changed: 111 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,111 @@
1+
package syncing
2+
3+
import (
4+
"context"
5+
"sync"
6+
7+
"github.com/evstack/ev-node/block/internal/common"
8+
"github.com/rs/zerolog"
9+
)
10+
11+
// AsyncDARetriever handles concurrent DA retrieval operations.
12+
type AsyncDARetriever struct {
13+
retriever DARetriever
14+
resultCh chan<- common.DAHeightEvent
15+
workCh chan uint64
16+
inFlight map[uint64]struct{}
17+
mu sync.Mutex
18+
logger zerolog.Logger
19+
wg sync.WaitGroup
20+
ctx context.Context
21+
cancel context.CancelFunc
22+
}
23+
24+
// NewAsyncDARetriever creates a new AsyncDARetriever.
25+
func NewAsyncDARetriever(
26+
retriever DARetriever,
27+
resultCh chan<- common.DAHeightEvent,
28+
logger zerolog.Logger,
29+
) *AsyncDARetriever {
30+
return &AsyncDARetriever{
31+
retriever: retriever,
32+
resultCh: resultCh,
33+
workCh: make(chan uint64, 100), // Buffer size 100
34+
inFlight: make(map[uint64]struct{}),
35+
logger: logger.With().Str("component", "async_da_retriever").Logger(),
36+
}
37+
}
38+
39+
// Start starts the worker pool.
40+
func (r *AsyncDARetriever) Start(ctx context.Context) {
41+
r.ctx, r.cancel = context.WithCancel(ctx)
42+
// Start 5 workers
43+
for i := 0; i < 5; i++ {
44+
r.wg.Add(1)
45+
go r.worker()
46+
}
47+
r.logger.Info().Msg("AsyncDARetriever started")
48+
}
49+
50+
// Stop stops the worker pool.
51+
func (r *AsyncDARetriever) Stop() {
52+
if r.cancel != nil {
53+
r.cancel()
54+
}
55+
r.wg.Wait()
56+
r.logger.Info().Msg("AsyncDARetriever stopped")
57+
}
58+
59+
// RequestRetrieval requests a DA retrieval for the given height.
60+
// It is non-blocking and idempotent.
61+
func (r *AsyncDARetriever) RequestRetrieval(height uint64) {
62+
r.mu.Lock()
63+
defer r.mu.Unlock()
64+
65+
if _, exists := r.inFlight[height]; exists {
66+
return
67+
}
68+
69+
select {
70+
case r.workCh <- height:
71+
r.inFlight[height] = struct{}{}
72+
r.logger.Debug().Uint64("height", height).Msg("queued DA retrieval request")
73+
default:
74+
r.logger.Debug().Uint64("height", height).Msg("DA retrieval worker pool full, dropping request")
75+
}
76+
}
77+
78+
func (r *AsyncDARetriever) worker() {
79+
defer r.wg.Done()
80+
81+
for {
82+
select {
83+
case <-r.ctx.Done():
84+
return
85+
case height := <-r.workCh:
86+
r.processRetrieval(height)
87+
}
88+
}
89+
}
90+
91+
func (r *AsyncDARetriever) processRetrieval(height uint64) {
92+
defer func() {
93+
r.mu.Lock()
94+
delete(r.inFlight, height)
95+
r.mu.Unlock()
96+
}()
97+
98+
events, err := r.retriever.RetrieveFromDA(r.ctx, height)
99+
if err != nil {
100+
r.logger.Debug().Err(err).Uint64("height", height).Msg("async DA retrieval failed")
101+
return
102+
}
103+
104+
for _, event := range events {
105+
select {
106+
case r.resultCh <- event:
107+
case <-r.ctx.Done():
108+
return
109+
}
110+
}
111+
}
Lines changed: 141 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,141 @@
1+
package syncing
2+
3+
import (
4+
"context"
5+
"testing"
6+
"time"
7+
8+
"github.com/evstack/ev-node/block/internal/common"
9+
"github.com/rs/zerolog"
10+
"github.com/stretchr/testify/assert"
11+
"github.com/stretchr/testify/mock"
12+
)
13+
14+
func TestAsyncDARetriever_RequestRetrieval(t *testing.T) {
15+
logger := zerolog.Nop()
16+
mockRetriever := NewMockDARetriever(t)
17+
resultCh := make(chan common.DAHeightEvent, 10)
18+
19+
asyncRetriever := NewAsyncDARetriever(mockRetriever, resultCh, logger)
20+
ctx, cancel := context.WithCancel(context.Background())
21+
defer cancel()
22+
23+
asyncRetriever.Start(ctx)
24+
defer asyncRetriever.Stop()
25+
26+
// 1. Test successful retrieval
27+
height1 := uint64(100)
28+
mockRetriever.EXPECT().RetrieveFromDA(mock.Anything, height1).Return([]common.DAHeightEvent{{DaHeight: height1}}, nil).Once()
29+
30+
asyncRetriever.RequestRetrieval(height1)
31+
32+
select {
33+
case event := <-resultCh:
34+
assert.Equal(t, height1, event.DaHeight)
35+
case <-time.After(1 * time.Second):
36+
t.Fatal("timeout waiting for result")
37+
}
38+
39+
// 2. Test deduplication (idempotency)
40+
// We'll block the retriever to simulate a slow request, then send multiple requests for the same height
41+
height2 := uint64(200)
42+
43+
// Create a channel to signal when the mock is called
44+
calledCh := make(chan struct{})
45+
// Create a channel to unblock the mock
46+
unblockCh := make(chan struct{})
47+
48+
mockRetriever.EXPECT().RetrieveFromDA(mock.Anything, height2).RunAndReturn(func(ctx context.Context, h uint64) ([]common.DAHeightEvent, error) {
49+
close(calledCh)
50+
<-unblockCh
51+
return []common.DAHeightEvent{{DaHeight: h}}, nil
52+
}).Once() // Should be called only once despite multiple requests
53+
54+
// Send first request
55+
asyncRetriever.RequestRetrieval(height2)
56+
57+
// Wait for the worker to pick it up
58+
select {
59+
case <-calledCh:
60+
case <-time.After(1 * time.Second):
61+
t.Fatal("timeout waiting for retriever call")
62+
}
63+
64+
// Send duplicate requests while the first one is still in flight
65+
asyncRetriever.RequestRetrieval(height2)
66+
asyncRetriever.RequestRetrieval(height2)
67+
68+
// Unblock the worker
69+
close(unblockCh)
70+
71+
// We should receive exactly one result
72+
select {
73+
case event := <-resultCh:
74+
assert.Equal(t, height2, event.DaHeight)
75+
case <-time.After(1 * time.Second):
76+
t.Fatal("timeout waiting for result")
77+
}
78+
79+
// Ensure no more results come through
80+
select {
81+
case <-resultCh:
82+
t.Fatal("received duplicate result")
83+
default:
84+
}
85+
}
86+
87+
func TestAsyncDARetriever_WorkerPoolLimit(t *testing.T) {
88+
logger := zerolog.Nop()
89+
mockRetriever := NewMockDARetriever(t)
90+
resultCh := make(chan common.DAHeightEvent, 100)
91+
92+
asyncRetriever := NewAsyncDARetriever(mockRetriever, resultCh, logger)
93+
ctx, cancel := context.WithCancel(context.Background())
94+
defer cancel()
95+
96+
asyncRetriever.Start(ctx)
97+
defer asyncRetriever.Stop()
98+
99+
// We have 5 workers. We'll block them all.
100+
unblockCh := make(chan struct{})
101+
102+
// Expect 5 calls that block
103+
for i := 0; i < 5; i++ {
104+
h := uint64(1000 + i)
105+
mockRetriever.EXPECT().RetrieveFromDA(mock.Anything, h).RunAndReturn(func(ctx context.Context, h uint64) ([]common.DAHeightEvent, error) {
106+
<-unblockCh
107+
return []common.DAHeightEvent{{DaHeight: h}}, nil
108+
}).Once()
109+
asyncRetriever.RequestRetrieval(h)
110+
}
111+
112+
// Give workers time to pick up tasks
113+
time.Sleep(100 * time.Millisecond)
114+
115+
// Now send a 6th request. It should be queued but not processed yet.
116+
height6 := uint64(1005)
117+
processed6 := make(chan struct{})
118+
mockRetriever.EXPECT().RetrieveFromDA(mock.Anything, height6).RunAndReturn(func(ctx context.Context, h uint64) ([]common.DAHeightEvent, error) {
119+
close(processed6)
120+
return []common.DAHeightEvent{{DaHeight: h}}, nil
121+
}).Once()
122+
123+
asyncRetriever.RequestRetrieval(height6)
124+
125+
// Ensure 6th request is NOT processed yet
126+
select {
127+
case <-processed6:
128+
t.Fatal("6th request processed too early")
129+
default:
130+
}
131+
132+
// Unblock workers
133+
close(unblockCh)
134+
135+
// Now 6th request should be processed
136+
select {
137+
case <-processed6:
138+
case <-time.After(1 * time.Second):
139+
t.Fatal("timeout waiting for 6th request")
140+
}
141+
}

block/internal/syncing/syncer.go

Lines changed: 11 additions & 23 deletions
Original file line numberDiff line numberDiff line change
@@ -68,6 +68,9 @@ type Syncer struct {
6868

6969
// P2P wait coordination
7070
p2pWaitState atomic.Value // stores p2pWaitState
71+
72+
// Async DA retriever
73+
asyncDARetriever *AsyncDARetriever
7174
}
7275

7376
// NewSyncer creates a new block syncer
@@ -114,6 +117,9 @@ func (s *Syncer) Start(ctx context.Context) error {
114117

115118
// Initialize handlers
116119
s.daRetriever = NewDARetriever(s.daClient, s.cache, s.genesis, s.logger)
120+
s.asyncDARetriever = NewAsyncDARetriever(s.daRetriever, s.heightInCh, s.logger)
121+
s.asyncDARetriever.Start(s.ctx)
122+
117123
s.p2pHandler = NewP2PHandler(s.headerStore, s.dataStore, s.cache, s.genesis, s.logger)
118124
if currentHeight, err := s.store.Height(s.ctx); err != nil {
119125
s.logger.Error().Err(err).Msg("failed to set initial processed height for p2p handler")
@@ -144,6 +150,9 @@ func (s *Syncer) Stop() error {
144150
if s.cancel != nil {
145151
s.cancel()
146152
}
153+
if s.asyncDARetriever != nil {
154+
s.asyncDARetriever.Stop()
155+
}
147156
s.cancelP2PWait(0)
148157
s.wg.Wait()
149158
s.logger.Info().Msg("syncer stopped")
@@ -466,29 +475,8 @@ func (s *Syncer) processHeightEvent(event *common.DAHeightEvent) {
466475
Uint64("da_height_hint", daHeightHint).
467476
Msg("P2P event with DA height hint, triggering targeted DA retrieval")
468477

469-
// Trigger targeted DA retrieval in background
470-
go func() {
471-
targetEvents, err := s.daRetriever.RetrieveFromDA(s.ctx, daHeightHint)
472-
if err != nil {
473-
s.logger.Debug().
474-
Err(err).
475-
Uint64("da_height", daHeightHint).
476-
Msg("targeted DA retrieval failed (hint may be incorrect or DA not yet available)")
477-
// Not a critical error - the sequential DA worker will eventually find it
478-
return
479-
}
480-
481-
// Process retrieved events from the targeted DA height
482-
for _, daEvent := range targetEvents {
483-
select {
484-
case s.heightInCh <- daEvent:
485-
case <-s.ctx.Done():
486-
return
487-
default:
488-
s.cache.SetPendingEvent(daEvent.Header.Height(), &daEvent)
489-
}
490-
}
491-
}()
478+
// Trigger targeted DA retrieval in background via worker pool
479+
s.asyncDARetriever.RequestRetrieval(daHeightHint)
492480
}
493481
}
494482
}

0 commit comments

Comments
 (0)