Skip to content

Commit 4e8d3ba

Browse files
committed
fix(da): fix polling fallback when ws not available
1 parent ec9f9bf commit 4e8d3ba

12 files changed

Lines changed: 248 additions & 11 deletions

File tree

‎CHANGELOG.md‎

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -9,6 +9,12 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0
99

1010
## [Unreleased]
1111

12+
## v1.1.4
13+
14+
### Fixed
15+
16+
- DA client falls back to HTTP polling with `Retrieve` when the WebSocket connection fails, instead of trying to use the WS-only `Subscribe` over HTTP. Automatically upgrade to WS is available [#3211](https://github.com/evstack/ev-node/pull/3361)
17+
1218
## v1.1.3
1319

1420
### Fixed

‎apps/evm/server/force_inclusion_test.go‎

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -85,6 +85,8 @@ func (m *mockDA) HasForcedInclusionNamespace() bool {
8585
return true
8686
}
8787

88+
func (m *mockDA) SupportsSubscribe() bool { return true }
89+
8890
func (m *mockDA) GetLatestDAHeight(_ context.Context) (uint64, error) {
8991
return 0, nil
9092
}

‎block/internal/da/async_block_retriever_test.go‎

Lines changed: 10 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -53,6 +53,8 @@ func TestAsyncBlockRetriever_SubscriptionDrivenCaching(t *testing.T) {
5353
client := &mocks.MockClient{}
5454
fiNs := datypes.NamespaceFromString("test-fi-ns").Bytes()
5555

56+
client.On("SupportsSubscribe").Return(true)
57+
5658
// Create a subscription channel that delivers one event then blocks.
5759
subCh := make(chan datypes.SubscriptionEvent, 1)
5860
subCh <- datypes.SubscriptionEvent{
@@ -104,6 +106,8 @@ func TestAsyncBlockRetriever_CatchupFillsGaps(t *testing.T) {
104106
client := &mocks.MockClient{}
105107
fiNs := datypes.NamespaceFromString("test-fi-ns").Bytes()
106108

109+
client.On("SupportsSubscribe").Return(true)
110+
107111
// Subscription delivers height 105 (no blobs — just a signal).
108112
subCh := make(chan datypes.SubscriptionEvent, 1)
109113
subCh <- datypes.SubscriptionEvent{Height: 105}
@@ -153,6 +157,8 @@ func TestAsyncBlockRetriever_HeightFromFuture(t *testing.T) {
153157
client := &mocks.MockClient{}
154158
fiNs := datypes.NamespaceFromString("test-fi-ns").Bytes()
155159

160+
client.On("SupportsSubscribe").Return(true)
161+
156162
// Subscription delivers height 100 with no blobs.
157163
subCh := make(chan datypes.SubscriptionEvent)
158164
client.On("Subscribe", mock.Anything, fiNs, mock.Anything).Return((<-chan datypes.SubscriptionEvent)(subCh), nil).Once()
@@ -187,6 +193,8 @@ func TestAsyncBlockRetriever_StopGracefully(t *testing.T) {
187193
client := &mocks.MockClient{}
188194
fiNs := datypes.NamespaceFromString("test-fi-ns").Bytes()
189195

196+
client.On("SupportsSubscribe").Return(true)
197+
190198
blockCh := make(chan datypes.SubscriptionEvent)
191199
client.On("Subscribe", mock.Anything, fiNs, mock.Anything).Return((<-chan datypes.SubscriptionEvent)(blockCh), nil).Maybe()
192200
client.On("Retrieve", mock.Anything, mock.Anything, fiNs).Return(datypes.ResultRetrieve{
@@ -211,6 +219,8 @@ func TestAsyncBlockRetriever_ReconnectOnSubscriptionError(t *testing.T) {
211219
client := &mocks.MockClient{}
212220
fiNs := datypes.NamespaceFromString("test-fi-ns").Bytes()
213221

222+
client.On("SupportsSubscribe").Return(true)
223+
214224
// First subscription closes immediately (simulating error).
215225
closedCh := make(chan datypes.SubscriptionEvent)
216226
close(closedCh)

‎block/internal/da/client.go‎

Lines changed: 8 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -38,6 +38,7 @@ type client struct {
3838
dataNamespaceBz []byte
3939
forcedNamespaceBz []byte
4040
hasForcedNamespace bool
41+
isWebSocket bool
4142
timestampCache *blockTimestampCache
4243
}
4344

@@ -137,6 +138,7 @@ func NewClient(cfg Config) FullClient {
137138
dataNamespaceBz: datypes.NamespaceFromString(cfg.DataNamespace).Bytes(),
138139
forcedNamespaceBz: forcedNamespaceBz,
139140
hasForcedNamespace: hasForcedNamespace,
141+
isWebSocket: cfg.DA.IsWebSocket,
140142
timestampCache: newBlockTimestampCache(blockTimestampCacheWindow),
141143
}
142144
}
@@ -485,6 +487,12 @@ func (c *client) HasForcedInclusionNamespace() bool {
485487
return c.hasForcedNamespace
486488
}
487489

490+
// SupportsSubscribe reports whether the underlying transport supports
491+
// channel-based subscriptions (WebSocket).
492+
func (c *client) SupportsSubscribe() bool {
493+
return c.isWebSocket
494+
}
495+
488496
// Subscribe subscribes to blobs in the given namespace via the celestia-node
489497
// Subscribe API. It returns a channel that emits a SubscriptionEvent for every
490498
// DA block containing a matching blob. The channel is closed when ctx is

‎block/internal/da/interface.go‎

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -31,6 +31,11 @@ type Client interface {
3131
// GetLatestDAHeight returns the latest height available on the DA layer.
3232
GetLatestDAHeight(ctx context.Context) (uint64, error)
3333

34+
// SupportsSubscribe reports whether the underlying transport supports
35+
// channel-based subscriptions (WebSocket). When false, callers must use
36+
// polling-based retrieval via Retrieve instead.
37+
SupportsSubscribe() bool
38+
3439
// Namespace accessors.
3540
GetHeaderNamespace() []byte
3641
GetDataNamespace() []byte

‎block/internal/da/subscriber.go‎

Lines changed: 69 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -122,10 +122,15 @@ func (s *Subscriber) Start(ctx context.Context) error {
122122

123123
ctx, cancel := context.WithCancel(ctx)
124124
s.cancel = cancel
125-
s.wg.Add(2)
126125
s.lifecycleMu.Unlock()
127126

128-
go s.followLoop(ctx)
127+
if s.client.SupportsSubscribe() {
128+
s.wg.Add(2)
129+
go s.followLoop(ctx)
130+
} else {
131+
s.wg.Add(2)
132+
go s.pollLoop(ctx)
133+
}
129134
go s.catchupLoop(ctx)
130135

131136
return nil
@@ -167,6 +172,68 @@ func (s *Subscriber) signalCatchup() {
167172
}
168173
}
169174

175+
// pollLoop periodically queries the latest DA height and triggers
176+
// catchup when new heights are available. The catchup loop fetches blobs
177+
// via Retrieve (which uses GetAll) so each height is fetched exactly once.
178+
// Periodically checks whether the underlying transport has been upgraded
179+
// to WebSocket and switches to followLoop when that happens.
180+
func (s *Subscriber) pollLoop(ctx context.Context) {
181+
defer s.wg.Done()
182+
183+
s.logger.Info().Msg("starting poll loop")
184+
defer s.logger.Info().Msg("poll loop stopped")
185+
186+
// Do an immediate poll on startup so we don't wait for the first tick.
187+
s.pollDAHeight(ctx)
188+
189+
ticker := time.NewTicker(s.daBlockTime)
190+
defer ticker.Stop()
191+
192+
for {
193+
// If the transport has been upgraded to WS in the background,
194+
// switch to the subscription-based follow loop.
195+
if s.client.SupportsSubscribe() {
196+
s.logger.Info().Msg("WebSocket available, switching from poll to follow loop")
197+
s.wg.Add(1)
198+
go s.followLoop(ctx)
199+
return
200+
}
201+
202+
select {
203+
case <-ctx.Done():
204+
return
205+
case <-ticker.C:
206+
s.pollDAHeight(ctx)
207+
}
208+
}
209+
}
210+
211+
// pollDAHeight queries GetLatestDAHeight and signals catchup when a new
212+
// height is observed. The actual blob retrieval is done by catchupLoop.
213+
func (s *Subscriber) pollDAHeight(ctx context.Context) {
214+
height, err := s.client.GetLatestDAHeight(ctx)
215+
if err != nil {
216+
if ctx.Err() != nil {
217+
return
218+
}
219+
s.logger.Warn().Err(err).Msg("poll: failed to get latest DA height")
220+
return
221+
}
222+
223+
cur := s.highestSeenDAHeight.Load()
224+
if height <= cur {
225+
return
226+
}
227+
228+
s.seenSubscriptionEvent.Store(true)
229+
s.logger.Debug().
230+
Uint64("new_da_height", height).
231+
Uint64("current_highest_seen", cur).
232+
Msg("poll: observed new DA height")
233+
234+
s.updateHighest(height)
235+
}
236+
170237
// followLoop subscribes to DA blob events and keeps highestSeenDAHeight up to date.
171238
func (s *Subscriber) followLoop(ctx context.Context) {
172239
defer s.wg.Done()

‎block/internal/da/tracing.go‎

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -165,6 +165,9 @@ func (t *tracedClient) GetForcedInclusionNamespace() []byte {
165165
func (t *tracedClient) HasForcedInclusionNamespace() bool {
166166
return t.inner.HasForcedInclusionNamespace()
167167
}
168+
func (t *tracedClient) SupportsSubscribe() bool {
169+
return t.inner.SupportsSubscribe()
170+
}
168171
func (t *tracedClient) Subscribe(ctx context.Context, namespace []byte, includeTimestamp bool) (<-chan datypes.SubscriptionEvent, error) {
169172
return t.inner.Subscribe(ctx, namespace, includeTimestamp)
170173
}

‎block/internal/da/tracing_test.go‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -74,6 +74,7 @@ func (m *mockFullClient) GetHeaderNamespace() []byte {
7474
func (m *mockFullClient) GetDataNamespace() []byte { return []byte{0x02} }
7575
func (m *mockFullClient) GetForcedInclusionNamespace() []byte { return []byte{0x03} }
7676
func (m *mockFullClient) HasForcedInclusionNamespace() bool { return true }
77+
func (m *mockFullClient) SupportsSubscribe() bool { return true }
7778

7879
// setup a tracer provider + span recorder
7980
func setupDATrace(t *testing.T, inner FullClient) (FullClient, *tracetest.SpanRecorder) {

‎block/internal/syncing/syncer_test.go‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -116,6 +116,7 @@ func makeSignedHeaderBytes(
116116
func setupMockDAClient(tb testing.TB) (da.Client, chan datypes.SubscriptionEvent) {
117117
mockClient := testmocks.NewMockClient(tb)
118118
eventCh := make(chan datypes.SubscriptionEvent, 1)
119+
mockClient.EXPECT().SupportsSubscribe().Return(true).Maybe()
119120
mockClient.EXPECT().Subscribe(mock.Anything, mock.Anything, mock.Anything).Return((<-chan datypes.SubscriptionEvent)(eventCh), nil).Maybe()
120121
return mockClient, eventCh
121122
}

‎pkg/da/jsonrpc/client.go‎

Lines changed: 94 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -5,6 +5,8 @@ import (
55
"fmt"
66
"net/http"
77
"strings"
8+
"sync"
9+
"time"
810

911
libshare "github.com/celestiaorg/go-square/v3/share"
1012
"github.com/filecoin-project/go-jsonrpc"
@@ -13,14 +15,28 @@ import (
1315

1416
// Client dials the celestia-node RPC "blob" and "header" namespaces.
1517
type Client struct {
16-
Blob BlobAPI
17-
Header HeaderAPI
18-
closer jsonrpc.ClientCloser
18+
Blob BlobAPI
19+
Header HeaderAPI
20+
IsWebSocket bool
21+
22+
mu sync.Mutex
23+
closer jsonrpc.ClientCloser
24+
retryCancel context.CancelFunc // stops the background WS retry loop
1925
}
2026

21-
// Close closes the underlying JSON-RPC connection.
27+
// Close closes the underlying JSON-RPC connection and stops any
28+
// background WebSocket retry loop.
2229
func (c *Client) Close() {
23-
if c != nil && c.closer != nil {
30+
if c == nil {
31+
return
32+
}
33+
c.mu.Lock()
34+
if c.retryCancel != nil {
35+
c.retryCancel()
36+
c.retryCancel = nil
37+
}
38+
c.mu.Unlock()
39+
if c.closer != nil {
2440
c.closer()
2541
}
2642
}
@@ -72,18 +88,88 @@ func NewClient(ctx context.Context, addr, token string, authHeaderName string) (
7288
// NewWSClient connects to the DA RPC endpoint over WebSocket.
7389
// Automatically converts http:// to ws:// (and https:// to wss://).
7490
// Supports channel-based subscriptions (e.g. Subscribe).
75-
// Note: WebSocket connections are eager — they connect at creation time
76-
// if the initial WS dial fails, falls back to HTTP polling for the entire session.
91+
// WebSocket connections are eager — they connect at creation time.
92+
// If the initial WS dial fails, it falls back to HTTP polling and spawns a
93+
// background goroutine that periodically retries the WS connection. When
94+
// the WS endpoint becomes reachable, the transport is transparently upgraded.
7795
func NewWSClient(ctx context.Context, logger zerolog.Logger, addr, token string, authHeaderName string) (*Client, error) {
7896
client, err := NewClient(ctx, httpToWS(addr), token, authHeaderName)
7997
if err != nil {
8098
logger.Warn().Err(err).Msg("DA websocket connection failed, falling back to DA polling")
81-
return NewClient(ctx, addr, token, authHeaderName)
99+
client, err = NewClient(ctx, addr, token, authHeaderName)
100+
if err != nil {
101+
return nil, err
102+
}
103+
client.IsWebSocket = false
104+
105+
// Retry WS in the background so transient outages don't force a permanent downgrade.
106+
retryCtx, retryCancel := context.WithCancel(context.Background())
107+
client.retryCancel = retryCancel
108+
go client.retryWSLoop(retryCtx, logger, addr, token, authHeaderName)
109+
110+
return client, nil
82111
}
83112

113+
client.IsWebSocket = true
84114
return client, nil
85115
}
86116

117+
const wsRetryInterval = 30 * time.Second
118+
119+
// retryWSLoop periodically attempts to re-establish a WebSocket connection.
120+
// When successful, it swaps the transport in-place and exits.
121+
func (c *Client) retryWSLoop(ctx context.Context, logger zerolog.Logger, addr, token, authHeaderName string) {
122+
ticker := time.NewTicker(wsRetryInterval)
123+
defer ticker.Stop()
124+
125+
for {
126+
select {
127+
case <-ctx.Done():
128+
return
129+
case <-ticker.C:
130+
if c.tryUpgradeWS(ctx, logger, addr, token, authHeaderName) {
131+
return
132+
}
133+
}
134+
}
135+
}
136+
137+
// tryUpgradeWS attempts to open a WS connection and, if successful, swaps
138+
// the transport internals so subsequent calls use WebSocket. Returns true
139+
// when the upgrade succeeds (or the client is already on WS).
140+
func (c *Client) tryUpgradeWS(ctx context.Context, logger zerolog.Logger, addr, token, authHeaderName string) bool {
141+
wsClient, err := NewClient(ctx, httpToWS(addr), token, authHeaderName)
142+
if err != nil {
143+
return false
144+
}
145+
146+
c.mu.Lock()
147+
defer c.mu.Unlock()
148+
149+
// Another goroutine may have already upgraded.
150+
if c.IsWebSocket {
151+
wsClient.Close()
152+
return true
153+
}
154+
155+
// Swap function pointers from the new WS client into the active client.
156+
c.Blob.Internal = wsClient.Blob.Internal
157+
c.Header.Internal = wsClient.Header.Internal
158+
159+
// Close the old HTTP connections and wire the new closer.
160+
oldCloser := c.closer
161+
c.closer = func() {
162+
wsClient.closer()
163+
if oldCloser != nil {
164+
oldCloser()
165+
}
166+
}
167+
168+
c.IsWebSocket = true
169+
logger.Info().Msg("DA websocket connection restored, switching back from HTTP polling")
170+
return true
171+
}
172+
87173
// BlobAPI mirrors celestia-node's blob module (nodebuilder/blob/blob.go).
88174
// jsonrpc.NewClient wires Internal.* to RPC stubs.
89175
type BlobAPI struct {

0 commit comments

Comments
 (0)