Skip to content

Commit caec54e

Browse files
committed
upgrade to ws
1 parent 8a31c2d commit caec54e

2 files changed

Lines changed: 96 additions & 6 deletions

File tree

‎block/internal/da/subscriber.go‎

Lines changed: 13 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -173,8 +173,10 @@ func (s *Subscriber) signalCatchup() {
173173
}
174174

175175
// pollLoop periodically queries the latest DA height and triggers
176-
// catchup when new heights are available. This is used when the
177-
// underlying transport does not support channel-based subscriptions (HTTP).
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.
178180
func (s *Subscriber) pollLoop(ctx context.Context) {
179181
defer s.wg.Done()
180182

@@ -188,6 +190,15 @@ func (s *Subscriber) pollLoop(ctx context.Context) {
188190
defer ticker.Stop()
189191

190192
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+
191202
select {
192203
case <-ctx.Done():
193204
return

‎pkg/da/jsonrpc/client.go‎

Lines changed: 83 additions & 4 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"
@@ -16,12 +18,25 @@ type Client struct {
1618
Blob BlobAPI
1719
Header HeaderAPI
1820
IsWebSocket bool
21+
22+
mu sync.Mutex
1923
closer jsonrpc.ClientCloser
24+
retryCancel context.CancelFunc // stops the background WS retry loop
2025
}
2126

22-
// Close closes the underlying JSON-RPC connection.
27+
// Close closes the underlying JSON-RPC connection and stops any
28+
// background WebSocket retry loop.
2329
func (c *Client) Close() {
24-
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 {
2540
c.closer()
2641
}
2742
}
@@ -73,8 +88,10 @@ func NewClient(ctx context.Context, addr, token string, authHeaderName string) (
7388
// NewWSClient connects to the DA RPC endpoint over WebSocket.
7489
// Automatically converts http:// to ws:// (and https:// to wss://).
7590
// Supports channel-based subscriptions (e.g. Subscribe).
76-
// Note: WebSocket connections are eager — they connect at creation time
77-
// 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.
7895
func NewWSClient(ctx context.Context, logger zerolog.Logger, addr, token string, authHeaderName string) (*Client, error) {
7996
client, err := NewClient(ctx, httpToWS(addr), token, authHeaderName)
8097
if err != nil {
@@ -84,13 +101,75 @@ func NewWSClient(ctx context.Context, logger zerolog.Logger, addr, token string,
84101
return nil, err
85102
}
86103
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+
87110
return client, nil
88111
}
89112

90113
client.IsWebSocket = true
91114
return client, nil
92115
}
93116

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+
94173
// BlobAPI mirrors celestia-node's blob module (nodebuilder/blob/blob.go).
95174
// jsonrpc.NewClient wires Internal.* to RPC stubs.
96175
type BlobAPI struct {

0 commit comments

Comments
 (0)