Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 2 additions & 2 deletions docs/pages/operation/websocket.mdx
Original file line number Diff line number Diff line change
Expand Up @@ -137,7 +137,7 @@ ws.send(JSON.stringify({

### Subscription features

- **Resilient delivery** — `newHeads` is subscribed on every connected WebSocket upstream of the network, and each block is delivered from whichever upstream announces it first. `logs` and `newPendingTransactions` are subscribed on every WebSocket upstream whose circuit breaker is closed; with [`failover.onDefaultsExhausted`](/config/projects/networks#failover) enabled, upstreams tagged `tier:fallback` are used only when every other upstream rejects the subscribe.
- **Resilient delivery** — `newHeads` is subscribed on every connected WebSocket upstream of the network, and each block is delivered from whichever upstream announces it first. With [`failover.onDefaultsExhausted`](/config/projects/networks#failover) enabled, heads from upstreams tagged `tier:fallback` are delivered only while no other upstream that the selection policy routes to has a live `newHeads` subscription, so clients are not told of a block only a fallback has while the primaries are up. A held-back head still updates that fallback's tip tracker. How quickly fallback heads take over when the primaries fail follows the selection policy's `evalInterval`. `logs` and `newPendingTransactions` are subscribed on every WebSocket upstream whose circuit breaker is closed; with [`failover.onDefaultsExhausted`](/config/projects/networks#failover) enabled, upstreams tagged `tier:fallback` are used only when every other upstream rejects the subscribe.
- **No surprise disconnects** — upstream reconnects, scoring changes and brief outages never close your connection; eRPC re-subscribes upstream behind the scenes. eRPC closes a connection only on shutdown (`1001`), when the client stops answering pings, or when a `logs` subscription overflows its buffer (see below).
- **Deduplication** — even though multiple upstreams are subscribed in parallel, you only receive each block, log or pending tx once. Logs are identified by `blockHash`, `transactionHash` and `logIndex`; a log missing any of them is passed through without deduplication.
- **Shared subscriptions** — clients subscribing to the same event on the same network share one upstream subscription per upstream. Each client gets its own subscription ID and receives every notification independently.
Expand Down Expand Up @@ -329,7 +329,7 @@ eRPC minimises the gap in three ways:

- **Per-block cache keying for moving tags.** The realtime cache entry for a `latest` / `finalized` request is keyed by the block number the tag currently resolves to, so each tip advance is a distinct key rather than pinning every call within the TTL to the same cached payload.
- **`enforceHighestBlock` response rewrite.** When the upstream returns a block number below the network's tracked tip, eRPC's [integrity check](/config/failsafe/integrity) fetches the specific higher block from another upstream and returns that instead (falling back to the highest valid block it has if none can serve it).
- **Delivered-head floor.** When an upstream WebSocket delivers a `newHeads` notification, eRPC updates that upstream's tip tracker before fanning the head out to clients, and records it as a floor for the network's `latest`. Every head a subscriber can receive feeds the floor once the delivering upstream's tip tracker has accepted it, and the floor is bounded: a value no live upstream comes close to (a stale or poisoned entry) is ignored rather than trusted, so it cannot pin the served tip. It never lifts `latest` past [`servedTip.guaranteedFor`](/config/projects/networks) / `guaranteedMethods`, and it does not apply to requests scoped with `use-upstream`. With shared state configured the floor is also propagated to other instances in the background, so they converge within the propagation delay rather than instantly.
- **Delivered-head floor.** When an upstream WebSocket delivers a `newHeads` notification, eRPC updates that upstream's tip tracker before fanning the head out to clients, and records it as a floor for the network's `latest`. Every head a subscriber can receive (so not a held-back fallback head) feeds the floor once the delivering upstream's tip tracker has accepted it, and the floor is bounded: a value no live upstream comes close to (a stale or poisoned entry) is ignored rather than trusted, so it cannot pin the served tip. It never lifts `latest` past [`servedTip.guaranteedFor`](/config/projects/networks) / `guaranteedMethods`, and it does not apply to requests scoped with `use-upstream`. With shared state configured the floor is also propagated to other instances in the background, so they converge within the propagation delay rather than instantly.

Even so, a transient single-block lag can reach a client when the chain produces a block faster than an in-flight HTTP round-trip completes. If your consumer treats single-block HTTP-vs-WS lag as a hard failure, widen its tolerance to match the chain's block time rather than expect eRPC (or any proxy) to eliminate it.

Expand Down
15 changes: 0 additions & 15 deletions erpc/network_executor.go
Original file line number Diff line number Diff line change
Expand Up @@ -656,21 +656,6 @@ func (e *networkExecutor) runHedge(
if uxe.Upstreams() == nil || len(uxe.Upstreams()) == 0 {
return false
}
// When every upstream returned -32004/missing-data, no sibling
// hedge leg can do better — they will all find the same upstreams
// consumed. Keep this result so the retry layer (shouldRetryWithReason)
// sees ErrUpstreamsExhausted directly and applies the 500ms delay.
causes := uxe.Errors()
allMissing := len(causes) > 0
for _, c := range causes {
if !common.HasErrorCode(c, common.ErrCodeEndpointMissingData) {
allMissing = false
break
}
}
if allMissing {
return true
}
}
// Underlying-retryable wrapped errors (e.g. ErrUpstreamsExhausted
// wrapping a 5xx) should continue racing for a healthier sibling.
Expand Down
9 changes: 6 additions & 3 deletions erpc/networks.go
Original file line number Diff line number Diff line change
Expand Up @@ -3627,9 +3627,12 @@ func (n *Network) acquireRateLimitPermit(ctx context.Context, req *common.Normal
// tierUpstreamsByGroup moves fallback-tier upstreams behind the rest,
// preserving order within each tier.
func tierUpstreamsByGroup(ups []common.Upstream) []common.Upstream {
return stablePartition(ups, func(u common.Upstream) bool {
return u.Config() != nil && u.Config().HasTag(common.TagTierFallback)
})
return stablePartition(ups, isFallbackTier)
}

// isFallbackTier reports whether u is tagged tier:fallback.
func isFallbackTier(u common.Upstream) bool {
return u.Config() != nil && u.Config().HasTag(common.TagTierFallback)
}

// partitionUpstreamsByLatestBlock moves upstreams whose polled head is known
Expand Down
13 changes: 12 additions & 1 deletion erpc/networks_failover_escape_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -169,6 +169,9 @@ func buildFailoverNetwork(
networkConfig.Failover = &common.FailoverConfig{OnDefaultsExhausted: util.BoolPtr(true)}
}
networkConfig.Failsafe = opts.failsafe
if opts.network != nil {
opts.network(networkConfig)
}

var policyEngine *policy.Engine
if !opts.noPolicy {
Expand Down Expand Up @@ -208,6 +211,10 @@ type failoverFixtureOpts struct {
// mocks registers test-specific mocks ahead of the standard ones, before
// any poller starts.
mocks func()
// configure adjusts the upstream configs before the network is built.
configure func(cfgs []*common.UpstreamConfig)
// network adjusts the network config before the network is built.
network func(cfg *common.NetworkConfig)
}

func setupFailoverFixture(
Expand All @@ -232,7 +239,11 @@ func setupFailoverFixture(
mockEthCallReturning("rpc3.localhost", "0x3333")
mockEthCallReturning("rpc4.localhost", "0x4444")

network, upr, mt := buildFailoverNetwork(t, ctx, failoverUpstreamConfigs(), opts)
cfgs := failoverUpstreamConfigs()
if opts.configure != nil {
opts.configure(cfgs)
}
network, upr, mt := buildFailoverNetwork(t, ctx, cfgs, opts)

upsList := upr.GetNetworkUpstreams(ctx, util.EvmNetworkId(999))
require.Len(t, upsList, 4)
Expand Down
167 changes: 167 additions & 0 deletions erpc/networks_failover_hedge_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,167 @@
package erpc

import (
"context"
"net/http"
"strings"
"sync/atomic"
"testing"
"time"

"github.com/erpc/erpc/common"
"github.com/erpc/erpc/telemetry"
"github.com/erpc/erpc/util"
"github.com/h2non/gock"
promUtil "github.com/prometheus/client_golang/prometheus/testutil"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
)

// unboundedAll drops the availability bound from every upstream, so each is
// called (and answers) for blocks above its polled head.
func unboundedAll(cfgs []*common.UpstreamConfig) {
for _, cfg := range cfgs {
cfg.Evm.BlockAvailability = nil
}
}

// ethCallMock answers eth_call on host with body after delay.
func ethCallMock(host string, delay time.Duration, body string) {
gock.New("http://" + host).
Post("").
Persist().
Filter(func(r *http.Request) bool {
return r.URL.Host == host && strings.Contains(util.SafeReadBody(r), "eth_call")
}).
Reply(200).
Delay(delay).
JSON([]byte(body))
}

// countEthCalls counts eth_call requests reaching host, ahead of its standard
// mock.
func countEthCalls(host string, hits *atomic.Int64) {
gock.New("http://" + host).
Post("").
Persist().
Filter(func(r *http.Request) bool {
// Filters run before host matching.
if r.URL.Host == host && strings.Contains(util.SafeReadBody(r), "eth_call") {
hits.Add(1)
}
return false
}).
Reply(200)
}

const missingDataBody = `{"jsonrpc":"2.0","id":1,"error":{"code":-32000,"message":"header not found"}}`

// hedgeFaster hedges well before a slow upstream answers.
func hedgeFaster() []*common.FailsafeConfig {
return []*common.FailsafeConfig{{
MatchMethod: "*",
Hedge: &common.HedgePolicyConfig{Delay: common.NewStaticDuration(20 * time.Millisecond), MaxCount: 1},
}}
}

func forwardEthCall(t *testing.T, ctx context.Context, network *Network, id int, blockHex string) (string, error) {
t.Helper()
req := ethCallRequest(id, blockHex)
req.SetNetwork(network)
resp, err := network.Forward(ctx, req)
if err != nil {
return "", err
}
require.NotNil(t, resp)
defer resp.Release()
jrr, err := resp.JsonRpcResponse()
require.NoError(t, err)
return strings.Trim(jrr.GetResultString(), `"`), nil
}

// A hedge leg that finds every routed upstream missing the data must not
// cancel a sibling leg that escalated to the fallbacks and is still waiting.
func TestFailover_HedgeKeepsEscalatedSibling(t *testing.T) {
defer util.ResetGock()
ctx, cancel := context.WithCancel(context.Background())
defer cancel()

// The primaries miss, the escape sends one leg to the (slow) fallbacks,
// and the hedge leg misses on the primaries again.
network, _, _ := setupFailoverFixture(t, ctx, failoverFixtureOpts{
primaryLatest: "0x3e8", // 1000
fallbackLatest: "0x3e8", // 1000
enableFailover: true,
configure: unboundedAll,
failsafe: hedgeFaster(),
mocks: func() {
ethCallMock("rpc1.localhost", 0, missingDataBody)
ethCallMock("rpc2.localhost", 0, missingDataBody)
ethCallMock("rpc3.localhost", 80*time.Millisecond, `{"jsonrpc":"2.0","id":1,"result":"0x3333"}`)
ethCallMock("rpc4.localhost", 80*time.Millisecond, `{"jsonrpc":"2.0","id":1,"result":"0x3333"}`)
},
})

escape := telemetry.MetricNetworkFallbackEscapeTotal.WithLabelValues("main", "evm:999", "eth_call")
escapeBefore := promUtil.ToFloat64(escape)
for i := 0; i < 5; i++ {
result, err := forwardEthCall(t, ctx, network, i, "0x3ea")
require.NoError(t, err, "iter %d", i)
assert.Equal(t, "0x3333", result, "iter %d", i)
}
assert.Equal(t, escapeBefore+5, promUtil.ToFloat64(escape))
}

// A fallback that already answered missing data must not let a hedge leg
// cancel another fallback that is still working on the request.
func TestFailover_HedgeKeepsSlowFallbackAfterFastFallbackMiss(t *testing.T) {
defer util.ResetGock()
ctx, cancel := context.WithCancel(context.Background())
defer cancel()

network, _, _ := setupFailoverFixture(t, ctx, failoverFixtureOpts{
primaryLatest: "0x3e8", // 1000
fallbackLatest: "0x3e8", // 1000
enableFailover: true,
configure: unboundedAll,
failsafe: hedgeFaster(),
mocks: func() {
ethCallMock("rpc1.localhost", 0, missingDataBody)
ethCallMock("rpc2.localhost", 0, missingDataBody)
ethCallMock("rpc3.localhost", 0, missingDataBody)
ethCallMock("rpc4.localhost", 80*time.Millisecond, `{"jsonrpc":"2.0","id":1,"result":"0x4444"}`)
},
})

for i := 0; i < 5; i++ {
result, err := forwardEthCall(t, ctx, network, i, "0x3ea")
require.NoError(t, err, "iter %d", i)
assert.Equal(t, "0x4444", result, "iter %d", i)
}
}

// A fallback whose polled head is ahead does not take reads the primaries
// serve: fallbacks are reached only through the escape. (The default
// policy's probeExcluded may still mirror the request to them in the
// background, so upstream hits are not counted here.)
func TestFailover_FallbackAheadDoesNotTakeServedReads(t *testing.T) {
defer util.ResetGock()
ctx, cancel := context.WithCancel(context.Background())
defer cancel()

network, _, _ := setupFailoverFixture(t, ctx, failoverFixtureOpts{
primaryLatest: "0x3e8", // 1000
fallbackLatest: "0x3ea", // 1002
enableFailover: true,
configure: unboundedAll,
})

escape := telemetry.MetricNetworkFallbackEscapeTotal.WithLabelValues("main", "evm:999", "eth_call")
escapeBefore := promUtil.ToFloat64(escape)
for i := 0; i < 5; i++ {
result, err := forwardEthCall(t, ctx, network, i, "0x3ea")
require.NoError(t, err, "iter %d", i)
assert.Contains(t, []string{"0x1111", "0x2222"}, result, "iter %d", i)
}
assert.Equal(t, escapeBefore, promUtil.ToFloat64(escape))
}
99 changes: 99 additions & 0 deletions erpc/networks_ws_tip_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@ import (
"github.com/erpc/erpc/common"
"github.com/erpc/erpc/data"
"github.com/erpc/erpc/health"
"github.com/erpc/erpc/internal/policy"
"github.com/erpc/erpc/thirdparty"
"github.com/erpc/erpc/upstream"
"github.com/erpc/erpc/util"
Expand Down Expand Up @@ -154,6 +155,104 @@ func TestNetworkHandle_SuggestLatestBlock_EveryKnownSourceLiftsLatest(t *testing
"a head a client can receive from the cordoned fallback floors latest")
}

type fakeHeadSource bool

func (f fakeHeadSource) HeadsLive() bool { return bool(f) }

// With failover on, a fallback's head reaches clients (and lifts latest) only
// when no primary the policy routes to is streaming heads itself. The
// fallback's own poller sees every head regardless.
func TestNetworkHandle_FallbackHeadsHeldWhilePrimariesStream(t *testing.T) {
setup := func(t *testing.T, ctx context.Context) (*Network, map[string]*upstream.Upstream) {
network, ups, _ := setupFailoverFixture(t, ctx, failoverFixtureOpts{
primaryLatest: "0x3e8", // 1000
fallbackLatest: "0x3e8",
enableFailover: true,
})
require.Equal(t, int64(1000), network.EvmHighestLatestBlockNumber(ctx))
byId := make(map[string]*upstream.Upstream, len(ups))
for _, u := range ups {
byId[u.Id()] = u
}
return network, byId
}

t.Run("HeldWhileAPrimaryStreams", func(t *testing.T) {
defer util.ResetGock()
ctx, cancel := context.WithCancel(context.Background())
defer cancel()
network, ups := setup(t, ctx)

handle := &networkHandle{nw: network, heads: map[string]headSource{
"primary-1": fakeHeadSource(false),
"primary-2": fakeHeadSource(true),
"fallback-1": fakeHeadSource(true),
}}
assert.False(t, handle.SuggestLatestBlock("ws:fallback-1", 1010))
assert.Equal(t, int64(1010), ups["fallback-1"].EvmStatePoller().LatestBlock(),
"the fallback's poller still sees its head")
assert.Equal(t, int64(1000), network.EvmHighestLatestBlockNumber(ctx),
"a held head must not lift latest")

assert.True(t, handle.SuggestLatestBlock("ws:primary-2", 1001))
assert.Equal(t, int64(1001), network.EvmHighestLatestBlockNumber(ctx))
})

t.Run("DeliveredWhenNoPrimaryStreams", func(t *testing.T) {
defer util.ResetGock()
ctx, cancel := context.WithCancel(context.Background())
defer cancel()
network, _ := setup(t, ctx)

// Primaries connected over HTTP only, or with no live subscription.
for _, heads := range []map[string]headSource{
{"fallback-1": fakeHeadSource(true)},
{"primary-1": fakeHeadSource(false), "fallback-1": fakeHeadSource(true)},
} {
handle := &networkHandle{nw: network, heads: heads}
assert.True(t, handle.SuggestLatestBlock("ws:fallback-1", 1010))
}
assert.Equal(t, int64(1010), network.EvmHighestLatestBlockNumber(ctx))
})

t.Run("DeliveredWhenPolicyExcludesPrimaries", func(t *testing.T) {
defer util.ResetGock()
ctx, cancel := context.WithCancel(context.Background())
defer cancel()
network, ups := setup(t, ctx)

handle := &networkHandle{nw: network, heads: map[string]headSource{
"primary-1": fakeHeadSource(true),
"primary-2": fakeHeadSource(true),
"fallback-1": fakeHeadSource(true),
}}
require.False(t, handle.SuggestLatestBlock("ws:fallback-1", 1010))

ups["primary-1"].Cordon("*", "test")
ups["primary-2"].Cordon("*", "test")
policy.TickForTest(network.policyEngine, network.networkId, "*")

assert.True(t, handle.SuggestLatestBlock("ws:fallback-1", 1011))
assert.Equal(t, int64(1011), network.EvmHighestLatestBlockNumber(ctx))
})

t.Run("DeliveredWithFailoverOff", func(t *testing.T) {
defer util.ResetGock()
ctx, cancel := context.WithCancel(context.Background())
defer cancel()
network, _, _ := setupFailoverFixture(t, ctx, failoverFixtureOpts{
primaryLatest: "0x3e8",
fallbackLatest: "0x3e8",
})

handle := &networkHandle{nw: network, heads: map[string]headSource{
"primary-1": fakeHeadSource(true),
"fallback-1": fakeHeadSource(true),
}}
assert.True(t, handle.SuggestLatestBlock("ws:fallback-1", 1010))
})
}

// The delivered-head floor is network-wide: a use-upstream-scoped request is
// answered from its group's own head.
func TestDeliveredHeadFloor_SkipsSelectorScopedRequests(t *testing.T) {
Expand Down
Loading
Loading