From 8e98cb28b67cc881d67456371225dc0696f695f1 Mon Sep 17 00:00:00 2001 From: shpookas Date: Wed, 30 Sep 2026 09:04:53 +0200 Subject: [PATCH 1/4] fix(network): route block-pinned requests to fallbacks that have the block - Tip-leader routing: a request pinned to a block above every routed upstream's known head is tried first on the tier:fallback upstreams whose head already reached it (their newHeads keep their pollers current), then on the routed list. It takes the per-request escalation, so other hedge legs and retries do not escape; the sweep that took it still escapes once to the fallbacks it has not tried. No-op when a routed head is unknown or at the block, for consensus, and with failover off. Leaders must match the use-upstream selector and their enforced availability bounds. Counted in erpc_network_tip_leader_route_total. - Hedge keeper: once a request has escalated, an all-missing ErrUpstreamsExhausted is not kept, so a leg that only re-swept the routed upstreams cannot cancel a leg still waiting on a fallback. If every leg misses, the hedge returns the last result. - Future-block short-circuit (served-tip): skip the synthetic null when a reachable fallback already has the block, except for consensus. --- common/request.go | 6 + docs/pages/config/projects/networks.mdx | 2 + docs/pages/reference/metrics.mdx | 1 + erpc/network_executor.go | 5 +- erpc/networks.go | 110 ++++- erpc/networks_failover_escape_test.go | 32 +- erpc/networks_tip_leader_test.go | 544 ++++++++++++++++++++++++ telemetry/metrics.go | 9 + 8 files changed, 701 insertions(+), 8 deletions(-) create mode 100644 erpc/networks_tip_leader_test.go diff --git a/common/request.go b/common/request.go index a68358e14..092c63e83 100644 --- a/common/request.go +++ b/common/request.go @@ -1281,6 +1281,12 @@ func (r *NormalizedRequest) MarkEscalatedToFallbacks() bool { return r.escalatedToFallbacks.CompareAndSwap(false, true) } +// EscalatedToFallbacks reports whether the request has spent its fallback +// escalation. +func (r *NormalizedRequest) EscalatedToFallbacks() bool { + return r != nil && r.escalatedToFallbacks.Load() +} + // UserId returns the user ID from the user object, or "n/a" if not available func (r *NormalizedRequest) UserId() string { if r == nil { diff --git a/docs/pages/config/projects/networks.mdx b/docs/pages/config/projects/networks.mdx index 33e967ce3..b69dfde43 100644 --- a/docs/pages/config/projects/networks.mdx +++ b/docs/pages/config/projects/networks.mdx @@ -147,6 +147,8 @@ The escape: - fires at most once per request, and never for consensus requests; - is counted in `erpc_network_fallback_escape_total{project,network,category}` — expect zero in steady state. +With `onDefaultsExhausted` on, a request pinned to a block above every routed upstream's known head is sent first to the `tier:fallback` upstreams whose head has already reached that block (for example because their `newHeads` subscription announced it first), then to the routed upstreams. This uses the same once-per-request escalation: other hedge legs and retries of that request do not escape, while the sweep that routed to the leaders still escapes once, to the fallbacks it has not tried, if the leaders and the routed upstreams all fail (counted in `erpc_network_fallback_escape_total` as usual). It does nothing when any routed upstream's head is unknown or already at the block, and never applies to consensus requests. It is counted in `erpc_network_tip_leader_route_total{project,network,category}`. + The validation report warns when `onDefaultsExhausted` is enabled but no upstream is tagged `tier:fallback`. ## Agent reference diff --git a/docs/pages/reference/metrics.mdx b/docs/pages/reference/metrics.mdx index b6e97c800..9d0e830b8 100644 --- a/docs/pages/reference/metrics.mdx +++ b/docs/pages/reference/metrics.mdx @@ -312,6 +312,7 @@ All metric names carry the `erpc_` prefix. Full definitions: 0 { + if fe := n.getFailsafeExecutor(ctx, req); fe == nil || !fe.HasConsensus() { + return nil, false + } + } jrr, err := common.NewJsonRpcResponse(req.ID(), nil, nil) if err != nil { return nil, false @@ -2294,6 +2301,37 @@ func (n *Network) Forward(ctx context.Context, req *common.NormalizedRequest) (* return requestBlockNumber(ctx, effectiveReq) > 0 && !evm.EmptyResultBeyondConfidence(ctx, effectiveReq) } + // Tip-leader routing: a request pinned to a block above every routed + // upstream's head goes first to the fallbacks that already have it + // (a fallback's WS heads keep its poller current), then to the routed + // list. It takes the per-request escalation, so no other leg or + // attempt escapes; this sweep keeps its escape to the fallbacks it + // has not tried. + leaderRouted := false + if !oneUpstreamOnly && n.cfg.Failover.Enabled() && !failsafeExecutor.HasConsensus() { + if leaders := n.tipLeaderFallbacks(execSpanCtx, effectiveReq, method, upsList); len(leaders) > 0 && + effectiveReq.MarkEscalatedToFallbacks() { + leaderRouted = true + routed := effectiveReq.NextUpstream + maxLoopIterations += len(leaders) + nextUpstream = func() (common.Upstream, error) { + for len(leaders) > 0 { + fb := leaders[0] + leaders = leaders[1:] + if _, loaded := effectiveReq.ConsumedUpstreams.LoadOrStore(fb, true); !loaded { + return fb, nil + } + } + return routed() + } + + telemetry.MetricNetworkTipLeaderRouteTotal.WithLabelValues( + n.projectId, n.Label(), method, + ).Inc() + lg.Debug().Int("leaders", len(leaders)).Msg("routing to fallback upstreams ahead of the routed tip") + } + } + escalationLoop: for { for loopIteration := 0; loopIteration < maxLoopIterations; loopIteration++ { @@ -2492,7 +2530,8 @@ func (n *Network) Forward(ctx context.Context, req *common.NormalizedRequest) (* fallbacks = append(fallbacks, fb) } } - if len(fallbacks) > 0 && effectiveReq.MarkEscalatedToFallbacks() { + if len(fallbacks) > 0 && (leaderRouted || effectiveReq.MarkEscalatedToFallbacks()) { + leaderRouted = false if bestResp != nil { bestResp.Release() bestResp = nil @@ -3627,9 +3666,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 @@ -3690,6 +3732,66 @@ func preferTipLeaderForNearTipGetBlock(ups []common.Upstream, method string, bn return out } +// tipLeaderFallbacks returns the fallback-tier upstreams, outside the routed +// list and allowed by the request's upstream selector, whose head has reached +// the block req is pinned to, when every routed upstream's head is known and +// below it. Nil otherwise. +func (n *Network) tipLeaderFallbacks(ctx context.Context, req *common.NormalizedRequest, method string, routed []common.Upstream) []common.Upstream { + if n.Architecture() != common.ArchitectureEvm { + return nil + } + bn := requestBlockNumber(ctx, req) + if bn <= 0 { + return nil + } + routedIds := make(map[string]struct{}, len(routed)) + for _, u := range routed { + if lb := upstreamLatestBlock(u); lb <= 0 || lb >= bn { + return nil + } + routedIds[u.Id()] = struct{}{} + } + return n.fallbacksAtBlock(ctx, req, method, bn, routedIds) +} + +// fallbacksAtBlock returns the fallback-escape upstreams, not in skip and +// allowed by the request's upstream selector, whose head has reached bn and +// whose enforced availability bounds admit it. +func (n *Network) fallbacksAtBlock(ctx context.Context, req *common.NormalizedRequest, method string, bn int64, skip map[string]struct{}) []common.Upstream { + selector := "" + if d := req.Directives(); d != nil { + selector = d.UseUpstream + } + var out []common.Upstream + for _, fb := range n.upstreamsRegistry.GetFallbackEscapeUpstreams(ctx, n.networkId, method) { + if _, ok := skip[fb.Id()]; ok || upstreamLatestBlock(fb) < bn || !n.availabilityAdmits(fb, method, bn) { + continue + } + if selector != "" { + if match, err := common.UpstreamMatchesSelector(selector, fb); err != nil || !match { + continue + } + } + out = append(out, fb) + } + return out +} + +// availabilityAdmits reports whether u's block-availability bounds, where +// enforced for method, admit bn. The same bounds checkUpstreamBlockAvailability +// gates on, without its metrics. +func (n *Network) availabilityAdmits(u common.Upstream, method string, bn int64) bool { + if methodHasDedicatedRangeAvailabilityHook(method) || n.blockAvailabilityExplicitlyDisabled(method) { + return true + } + eu, ok := u.(common.EvmUpstream) + if !ok { + return true + } + lo, hi := eu.EvmBlockAvailabilityBounds() + return (lo == math.MinInt64 || bn >= lo) && (hi == math.MaxInt64 || bn <= hi) +} + // upstreamLatestBlock is u's polled head, or 0 when unknown. func upstreamLatestBlock(u common.Upstream) int64 { if eu, ok := u.(common.EvmUpstream); ok { diff --git a/erpc/networks_failover_escape_test.go b/erpc/networks_failover_escape_test.go index 1f9f60cac..945a59372 100644 --- a/erpc/networks_failover_escape_test.go +++ b/erpc/networks_failover_escape_test.go @@ -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 { @@ -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( @@ -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) @@ -265,11 +276,26 @@ func TestFailover_EscapeHatch(t *testing.T) { ctx, cancel := context.WithCancel(context.Background()) defer cancel() - // Primaries at 1000 skip block 1002; the fallbacks at 1002 serve it. + // Primaries at 1002 fail eth_call with a retryable error; the + // fallbacks at 1002 serve it. (Primaries below the block would take + // tip-leader routing instead, see TestFailover_TipLeaderRouting.) network, _, _ := setupFailoverFixture(t, ctx, failoverFixtureOpts{ - primaryLatest: "0x3e8", // 1000 + primaryLatest: "0x3ea", // 1002 fallbackLatest: "0x3ea", // 1002 enableFailover: true, + mocks: func() { + for _, host := range []string{"rpc1.localhost", "rpc2.localhost"} { + host := host + 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). + JSON([]byte(`{"jsonrpc":"2.0","id":1,"error":{"code":-32603,"message":"internal error"}}`)) + } + }, }) counter := telemetry.MetricNetworkFallbackEscapeTotal.WithLabelValues("main", "evm:999", "eth_call") diff --git a/erpc/networks_tip_leader_test.go b/erpc/networks_tip_leader_test.go new file mode 100644 index 000000000..d9241f83a --- /dev/null +++ b/erpc/networks_tip_leader_test.go @@ -0,0 +1,544 @@ +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" +) + +// unboundedPrimaries drops the availability bound from the primaries, so they +// are called (and answer) for blocks above their polled head, like a local +// node whose poller trails the head a fallback just announced. +func unboundedPrimaries(cfgs []*common.UpstreamConfig) { + for _, cfg := range cfgs { + if !cfg.HasTag(common.TagTierFallback) { + cfg.Evm.BlockAvailability = nil + } + } +} + +// 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) +} + +// missingOnPrimariesSlowOnFallbacks makes the primaries answer eth_call with +// missing data and the fallbacks serve it after delay. +func missingOnPrimariesSlowOnFallbacks(delay time.Duration) func() { + return func() { + for _, host := range []string{"rpc1.localhost", "rpc2.localhost"} { + host := host + 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). + JSON([]byte(`{"jsonrpc":"2.0","id":1,"error":{"code":-32000,"message":"header not found"}}`)) + } + for _, host := range []string{"rpc3.localhost", "rpc4.localhost"} { + host := host + 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(`{"jsonrpc":"2.0","id":1,"result":"0x3333"}`)) + } + } +} + +// hedgeFaster hedges well before a fallback 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 +} + +func TestFailover_TipLeaderRouting(t *testing.T) { + leaderCounter := func() float64 { + return promUtil.ToFloat64(telemetry.MetricNetworkTipLeaderRouteTotal.WithLabelValues("main", "evm:999", "eth_call")) + } + escapeCounter := func() float64 { + return promUtil.ToFloat64(telemetry.MetricNetworkFallbackEscapeTotal.WithLabelValues("main", "evm:999", "eth_call")) + } + + t.Run("RoutesToFallbackThatHasTheBlock", func(t *testing.T) { + defer util.ResetGock() + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + + // Primaries polled at 1000 would still answer for 1002; the fallbacks + // already have 1002, so they go first. + var primaryHits atomic.Int64 + network, _, _ := setupFailoverFixture(t, ctx, failoverFixtureOpts{ + primaryLatest: "0x3e8", // 1000 + fallbackLatest: "0x3ea", // 1002 + enableFailover: true, + configure: unboundedPrimaries, + mocks: func() { + countEthCalls("rpc1.localhost", &primaryHits) + countEthCalls("rpc2.localhost", &primaryHits) + }, + }) + + leaderBefore, escapeBefore := leaderCounter(), escapeCounter() + result, err := forwardEthCall(t, ctx, network, 1, "0x3ea") + require.NoError(t, err) + assert.Contains(t, []string{"0x3333", "0x4444"}, result, "a fallback that has the block must serve it") + assert.Equal(t, int64(0), primaryHits.Load(), "primaries must not be tried before the leader") + assert.Equal(t, leaderBefore+1, leaderCounter()) + assert.Equal(t, escapeBefore, escapeCounter(), "leader routing is not an escape") + }) + + t.Run("NoLeaderWhenRoutedHasTheBlock", func(t *testing.T) { + defer util.ResetGock() + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + + network, _, _ := setupFailoverFixture(t, ctx, failoverFixtureOpts{ + primaryLatest: "0x3ea", // 1002 + fallbackLatest: "0x3eb", // 1003 + enableFailover: true, + }) + + leaderBefore := leaderCounter() + for i := 0; i < 10; i++ { + result, err := forwardEthCall(t, ctx, network, i, "0x3ea") + require.NoError(t, err) + assert.Contains(t, []string{"0x1111", "0x2222"}, result, "a primary that has the block must serve it (iter %d)", i) + } + assert.Equal(t, leaderBefore, leaderCounter()) + }) + + t.Run("NoLeaderWhenFallbackIsBehindToo", func(t *testing.T) { + defer util.ResetGock() + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + + network, _, _ := setupFailoverFixture(t, ctx, failoverFixtureOpts{ + primaryLatest: "0x3e8", // 1000 + fallbackLatest: "0x3e9", // 1001 + enableFailover: true, + configure: unboundedPrimaries, + }) + + leaderBefore := leaderCounter() + result, err := forwardEthCall(t, ctx, network, 1, "0x3ea") + require.NoError(t, err) + assert.Contains(t, []string{"0x1111", "0x2222"}, result) + assert.Equal(t, leaderBefore, leaderCounter()) + }) + + t.Run("NoLeaderWhenFailoverDisabled", func(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: false, + configure: unboundedPrimaries, + }) + + leaderBefore := leaderCounter() + _, _ = forwardEthCall(t, ctx, network, 1, "0x3ea") + assert.Equal(t, leaderBefore, leaderCounter()) + }) + + t.Run("FallsBackToRoutedWhenLeaderFails", func(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: unboundedPrimaries, + mocks: func() { + for _, host := range []string{"rpc3.localhost", "rpc4.localhost"} { + host := host + 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). + JSON([]byte(`{"jsonrpc":"2.0","id":1,"error":{"code":-32603,"message":"internal error"}}`)) + } + }, + }) + + leaderBefore, escapeBefore := leaderCounter(), escapeCounter() + result, err := forwardEthCall(t, ctx, network, 1, "0x3ea") + require.NoError(t, err) + assert.Contains(t, []string{"0x1111", "0x2222"}, result, "routed upstreams must still be tried after the leaders fail") + assert.Equal(t, leaderBefore+1, leaderCounter()) + assert.Equal(t, escapeBefore, escapeCounter(), "the escalation is already spent") + }) + + t.Run("HedgeDoesNotCancelSlowLeader", func(t *testing.T) { + defer util.ResetGock() + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + + // The hedge leg sweeps the primaries (missing data) while the leader + // leg still waits on the fallback; the leader must win. + network, _, _ := setupFailoverFixture(t, ctx, failoverFixtureOpts{ + primaryLatest: "0x3e8", // 1000 + fallbackLatest: "0x3ea", // 1002 + enableFailover: true, + configure: unboundedPrimaries, + failsafe: hedgeFaster(), + mocks: missingOnPrimariesSlowOnFallbacks(80 * time.Millisecond), + }) + + 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) + } + }) +} + +// 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() + + // Nobody's head reaches 1002, so there is no tip leader: 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: func(cfgs []*common.UpstreamConfig) { + for _, cfg := range cfgs { + cfg.Evm.BlockAvailability = nil + } + }, + failsafe: hedgeFaster(), + mocks: missingOnPrimariesSlowOnFallbacks(80 * time.Millisecond), + }) + + escapeBefore := promUtil.ToFloat64(telemetry.MetricNetworkFallbackEscapeTotal.WithLabelValues("main", "evm:999", "eth_call")) + 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(telemetry.MetricNetworkFallbackEscapeTotal.WithLabelValues("main", "evm:999", "eth_call"))) +} + +// 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)) +} + +// latestHeadMock pins host's polled latest block, ahead of its standard mock. +func latestHeadMock(host, latestHex string) { + gock.New("http://" + host). + Post(""). + Persist(). + Filter(func(r *http.Request) bool { + b := util.SafeReadBody(r) + return r.URL.Host == host && strings.Contains(b, "eth_getBlockByNumber") && strings.Contains(b, `"latest"`) + }). + Reply(200). + JSON([]byte(`{"result":{"number":"` + latestHex + `","timestamp":"0x6702a8f0"}}`)) +} + +func unboundedAll(cfgs []*common.UpstreamConfig) { + for _, cfg := range cfgs { + cfg.Evm.BlockAvailability = nil + } +} + +const missingDataBody = `{"jsonrpc":"2.0","id":1,"error":{"code":-32000,"message":"header not found"}}` + +// 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, no tip leader + 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) + } +} + +// When the tip leader fails, the sweep still escapes to the fallbacks it has +// not tried, even one whose polled head trails the block. +func TestFailover_TipLeaderKeepsEscapeToOtherFallbacks(t *testing.T) { + defer util.ResetGock() + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + + network, _, _ := setupFailoverFixture(t, ctx, failoverFixtureOpts{ + primaryLatest: "0x3e8", // 1000 + fallbackLatest: "0x3ea", // fallback-1 at 1002 (leader) + enableFailover: true, + configure: unboundedAll, + mocks: func() { + latestHeadMock("rpc4.localhost", "0x3e8") // fallback-2 at 1000 + ethCallMock("rpc1.localhost", 0, missingDataBody) + ethCallMock("rpc2.localhost", 0, missingDataBody) + ethCallMock("rpc3.localhost", 0, `{"jsonrpc":"2.0","id":1,"error":{"code":-32603,"message":"internal error"}}`) + }, + }) + + leader := telemetry.MetricNetworkTipLeaderRouteTotal.WithLabelValues("main", "evm:999", "eth_call") + escape := telemetry.MetricNetworkFallbackEscapeTotal.WithLabelValues("main", "evm:999", "eth_call") + leaderBefore, escapeBefore := promUtil.ToFloat64(leader), promUtil.ToFloat64(escape) + + result, err := forwardEthCall(t, ctx, network, 1, "0x3ea") + require.NoError(t, err) + assert.Equal(t, "0x4444", result, "the untried fallback must still serve the request") + assert.Equal(t, leaderBefore+1, promUtil.ToFloat64(leader)) + assert.Equal(t, escapeBefore+1, promUtil.ToFloat64(escape)) +} + +// A request pinned to an upstream by the use-upstream directive never takes +// the tip-leader route to a fallback outside it. +func TestFailover_TipLeaderRespectsUseUpstream(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: unboundedPrimaries, + }) + + leader := telemetry.MetricNetworkTipLeaderRouteTotal.WithLabelValues("main", "evm:999", "eth_call") + before := promUtil.ToFloat64(leader) + + req := ethCallRequest(1, "0x3ea") + req.SetDirectives(&common.RequestDirectives{UseUpstream: "primary-1"}) + req.SetNetwork(network) + resp, err := network.Forward(ctx, req) + require.NoError(t, err) + require.NotNil(t, resp) + defer resp.Release() + jrr, err := resp.JsonRpcResponse() + require.NoError(t, err) + assert.Equal(t, "0x1111", strings.Trim(jrr.GetResultString(), `"`)) + assert.Equal(t, before, promUtil.ToFloat64(leader)) +} + +// With served-tip on, a numbered eth_getBlockByNumber above every eligible +// head is short-circuited to null, unless a reachable fallback already has +// the block. +func TestFailover_FutureBlockShortCircuitSparesFallbackThatHasIt(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, + network: func(cfg *common.NetworkConfig) { + cfg.Evm.ServedTip = &common.EvmServedTipConfig{EnabledFor: []string{"latest"}} + }, + mocks: func() { + for _, host := range []string{"rpc3.localhost", "rpc4.localhost"} { + host := host + gock.New("http://" + host). + Post(""). + Persist(). + Filter(func(r *http.Request) bool { + b := util.SafeReadBody(r) + return r.URL.Host == host && strings.Contains(b, "eth_getBlockByNumber") && strings.Contains(b, `"0x3ea"`) + }). + Reply(200). + JSON([]byte(`{"jsonrpc":"2.0","id":1,"result":{"number":"0x3ea","hash":"0xfb","timestamp":"0x6702a8f2"}}`)) + } + }, + }) + + req := common.NewNormalizedRequest([]byte(`{"jsonrpc":"2.0","id":1,"method":"eth_getBlockByNumber","params":["0x3ea",false]}`)) + req.SetNetwork(network) + resp, err := network.Forward(ctx, req) + require.NoError(t, err) + require.NotNil(t, resp) + defer resp.Release() + jrr, err := resp.JsonRpcResponse() + require.NoError(t, err) + assert.Contains(t, jrr.GetResultString(), `"0x3ea"`, "the fallback's block must be returned, not a synthetic null") +} + +// servedTipLatest turns served-tip on for the latest axis. +func servedTipLatest(cfg *common.NetworkConfig) { + cfg.Evm.ServedTip = &common.EvmServedTipConfig{EnabledFor: []string{"latest"}} +} + +func getBlockRequest(blockHex string) *common.NormalizedRequest { + return common.NewNormalizedRequest([]byte(`{"jsonrpc":"2.0","id":1,"method":"eth_getBlockByNumber","params":["` + blockHex + `",false]}`)) +} + +// Consensus never reaches a fallback, so a fallback that has the block must +// not stop the future-block short-circuit for it. +func TestFailover_FutureBlockShortCircuitStillAppliesToConsensus(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, + network: servedTipLatest, + failsafe: []*common.FailsafeConfig{{ + MatchMethod: "*", + Consensus: &common.ConsensusPolicyConfig{MaxParticipants: 2, AgreementThreshold: 2}, + }}, + }) + + req := getBlockRequest("0x3ea") + req.SetNetwork(network) + resp, err := network.Forward(ctx, req) + require.NoError(t, err, "consensus must still get the truthful null") + require.NotNil(t, resp) + defer resp.Release() + assert.True(t, resp.IsResultEmptyish(), "expected the short-circuit null") +} + +// A fallback whose configured availability excludes the block is neither a +// tip leader nor a reason to skip the future-block short-circuit. +func TestFailover_FallbackAvailabilityBoundsExcludeLeader(t *testing.T) { + leader := telemetry.MetricNetworkTipLeaderRouteTotal.WithLabelValues("main", "evm:999", "eth_call") + capFallbacks := func(cfgs []*common.UpstreamConfig) { + unboundedPrimaries(cfgs) + for _, cfg := range cfgs { + if cfg.HasTag(common.TagTierFallback) { + cfg.Evm.BlockAvailability = &common.EvmBlockAvailabilityConfig{ + Upper: &common.EvmAvailabilityBoundConfig{ExactBlock: i64(1000)}, + } + } + } + } + + t.Run("NotALeader", func(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, but capped at 1000 + enableFailover: true, + configure: capFallbacks, + }) + + before := promUtil.ToFloat64(leader) + result, err := forwardEthCall(t, ctx, network, 1, "0x3ea") + require.NoError(t, err) + assert.Contains(t, []string{"0x1111", "0x2222"}, result) + assert.Equal(t, before, promUtil.ToFloat64(leader)) + }) + + t.Run("ShortCircuitStillApplies", func(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, but capped at 1000 + enableFailover: true, + network: servedTipLatest, + configure: func(cfgs []*common.UpstreamConfig) { + capFallbacks(cfgs) + for _, cfg := range cfgs { + if !cfg.HasTag(common.TagTierFallback) { + cfg.Evm.BlockAvailability = headBoundedAvailability() + } + } + }, + }) + + req := getBlockRequest("0x3ea") + req.SetNetwork(network) + resp, err := network.Forward(ctx, req) + require.NoError(t, err) + require.NotNil(t, resp) + defer resp.Release() + assert.True(t, resp.IsResultEmptyish(), "expected the short-circuit null") + }) +} diff --git a/telemetry/metrics.go b/telemetry/metrics.go index 033c525e5..73fc3968e 100644 --- a/telemetry/metrics.go +++ b/telemetry/metrics.go @@ -759,6 +759,15 @@ var ( Help: "Total number of times the per-request fallback escape hatch fired because the primary upstream set was exhausted with retryable errors.", }, []string{"project", "network", "category"}) + // MetricNetworkTipLeaderRouteTotal counts block-pinned requests sent to a + // fallback-tier upstream first because it already had the block while + // every routed upstream's head was below it. + MetricNetworkTipLeaderRouteTotal = DefineCounter(prometheus.CounterOpts{ + Namespace: "erpc", + Name: "network_tip_leader_route_total", + Help: "Total number of block-pinned requests routed to a fallback-tier upstream first because it had the block and no routed upstream did.", + }, []string{"project", "network", "category"}) + // MetricCacheExecutorAttempt counts every attempt a cache-connector // failsafe executor governs, keyed by the executor's identity (the // matchMethod / matchFinality it was configured with) and the outcome From 5f4ecc1193be99e18b5739e8c33221656a3f59c9 Mon Sep 17 00:00:00 2001 From: Jonny Date: Mon, 5 Oct 2026 18:55:23 +0100 Subject: [PATCH 2/4] refactor(network): route tip leaders by ordering, not a second escalation Same outcome as the previous commit, with fewer moving parts: - Fallbacks whose polled head has the block join the request's upstream list, and tiering now runs before the has-the-block partition, so an upstream that has the block goes first whatever its tier. This is ordering, not an escalation: drops leaderRouted, the shared MarkEscalatedToFallbacks handling, the maxLoopIterations bump and NormalizedRequest.EscalatedToFallbacks. The once-per-request escape is unchanged; use-upstream is enforced by NextUpstream as for any upstream. - The future-block short-circuit is skipped when a tip leader was added, instead of rescanning the fallbacks. - The hedge keeper no longer keeps an all-missing ErrUpstreamsExhausted at all, rather than only after escalation: a sibling leg may be on an upstream this leg never tried, which already happens in the plain escape path with no tip leader (TestFailover_HedgeKeepsEscalatedSibling and TestFailover_HedgeKeepsSlowFallbackAfterFastFallbackMiss fail on 05126d01). If every leg ends that way the hedge still returns the last one to the retry layer. Tests unchanged; TestFailover_* pass, including under -race. Co-Authored-By: Claude Opus 5.5 --- common/request.go | 6 -- docs/pages/config/projects/networks.mdx | 2 +- erpc/network_executor.go | 18 ---- erpc/networks.go | 117 ++++++++---------------- 4 files changed, 41 insertions(+), 102 deletions(-) diff --git a/common/request.go b/common/request.go index 092c63e83..a68358e14 100644 --- a/common/request.go +++ b/common/request.go @@ -1281,12 +1281,6 @@ func (r *NormalizedRequest) MarkEscalatedToFallbacks() bool { return r.escalatedToFallbacks.CompareAndSwap(false, true) } -// EscalatedToFallbacks reports whether the request has spent its fallback -// escalation. -func (r *NormalizedRequest) EscalatedToFallbacks() bool { - return r != nil && r.escalatedToFallbacks.Load() -} - // UserId returns the user ID from the user object, or "n/a" if not available func (r *NormalizedRequest) UserId() string { if r == nil { diff --git a/docs/pages/config/projects/networks.mdx b/docs/pages/config/projects/networks.mdx index b69dfde43..352da637b 100644 --- a/docs/pages/config/projects/networks.mdx +++ b/docs/pages/config/projects/networks.mdx @@ -147,7 +147,7 @@ The escape: - fires at most once per request, and never for consensus requests; - is counted in `erpc_network_fallback_escape_total{project,network,category}` — expect zero in steady state. -With `onDefaultsExhausted` on, a request pinned to a block above every routed upstream's known head is sent first to the `tier:fallback` upstreams whose head has already reached that block (for example because their `newHeads` subscription announced it first), then to the routed upstreams. This uses the same once-per-request escalation: other hedge legs and retries of that request do not escape, while the sweep that routed to the leaders still escapes once, to the fallbacks it has not tried, if the leaders and the routed upstreams all fail (counted in `erpc_network_fallback_escape_total` as usual). It does nothing when any routed upstream's head is unknown or already at the block, and never applies to consensus requests. It is counted in `erpc_network_tip_leader_route_total{project,network,category}`. +With `onDefaultsExhausted` on, a request pinned to a block above every routed upstream's known head also routes to the `tier:fallback` upstreams whose head has already reached that block (for example because their `newHeads` subscription announced it first). For any block-pinned request, upstreams whose head has the block are tried before those whose head is below it, whatever their tier. This is ordering, not an escape: the escape above still fires once, to the fallbacks not yet tried, if every upstream fails. Fallbacks are added only when every routed upstream's head is known and below the block, and never for consensus requests; each addition is counted in `erpc_network_tip_leader_route_total{project,network,category}`. The validation report warns when `onDefaultsExhausted` is enabled but no upstream is tagged `tier:fallback`. diff --git a/erpc/network_executor.go b/erpc/network_executor.go index 57aecf8d6..8be2b55b2 100644 --- a/erpc/network_executor.go +++ b/erpc/network_executor.go @@ -656,24 +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 - } - } - // Unless the request escalated to the fallbacks: a sibling leg - // may still be waiting on one that has the data, so keep racing - // (if every leg ends like this the hedge returns the last one). - if allMissing { - return !req.EscalatedToFallbacks() - } } // Underlying-retryable wrapped errors (e.g. ErrUpstreamsExhausted // wrapping a 5xx) should continue racing for a healthier sibling. diff --git a/erpc/networks.go b/erpc/networks.go index f505ff458..cc9719a01 100644 --- a/erpc/networks.go +++ b/erpc/networks.go @@ -915,13 +915,6 @@ func (n *Network) tryShortCircuitFutureBlock(ctx context.Context, req *common.No // Unknown head (fail open) or block within reach of some upstream. return nil, false } - // A fallback the sweep can still reach may already have it while the - // policy keeps it out of the eligible set. Consensus never reaches it. - if n.cfg.Failover.Enabled() && len(n.fallbacksAtBlock(ctx, req, method, bn, nil)) > 0 { - if fe := n.getFailsafeExecutor(ctx, req); fe == nil || !fe.HasConsensus() { - return nil, false - } - } jrr, err := common.NewJsonRpcResponse(req.ID(), nil, nil) if err != nil { return nil, false @@ -2048,12 +2041,21 @@ func (n *Network) Forward(ctx context.Context, req *common.NormalizedRequest) (* return nil, err } - // For a specific block, try upstreams whose poller already has it first. - // Ordering hints only: every upstream stays eligible. + // For a specific block above every routed head, the fallbacks whose head + // already has it (e.g. their newHeads led) join the routed list. + var bn int64 + var tipLeaders []common.Upstream if n.Architecture() == common.ArchitectureEvm { - if bn := requestBlockNumber(ctx, req); bn > 0 { - upsList = partitionUpstreamsByLatestBlock(upsList, bn) - upsList = preferTipLeaderForNearTipGetBlock(upsList, method, bn) + bn = requestBlockNumber(ctx, req) + } + if bn > 0 && n.cfg.Failover.Enabled() { + if fe := n.getFailsafeExecutor(ctx, req); fe != nil && !fe.HasConsensus() { + if tipLeaders = n.tipLeaderFallbacks(ctx, req, method, bn, upsList); len(tipLeaders) > 0 { + upsList = append(slices.Clone(upsList), tipLeaders...) + telemetry.MetricNetworkTipLeaderRouteTotal.WithLabelValues( + n.projectId, n.Label(), method, + ).Inc() + } } } @@ -2063,6 +2065,13 @@ func (n *Network) Forward(ctx context.Context, req *common.NormalizedRequest) (* upsList = tierUpstreamsByGroup(upsList) } + // For a specific block, try upstreams whose poller already has it first, + // whatever their tier. Ordering hints only: every upstream stays eligible. + if bn > 0 { + upsList = partitionUpstreamsByLatestBlock(upsList, bn) + upsList = preferTipLeaderForNearTipGetBlock(upsList, method, bn) + } + // Architecture-specific pruning of the upstream list. Currently only SVM // uses this hook; both filters are gated on ArchitectureSvm so EVM networks // never enter this block. @@ -2160,13 +2169,16 @@ func (n *Network) Forward(ctx context.Context, req *common.NormalizedRequest) (* // Future-block short-circuit: a concrete block number beyond every eligible // upstream's head cannot be served yet — return the truthful null instead of // dispatching + hedging across upstreams that will all return empty. Runs - // before rate limiting so a non-dispatched request consumes no permit. - if resp, ok := n.tryShortCircuitFutureBlock(ctx, req, method); ok { - forwardSpan.SetAttributes(attribute.Bool("future_block.short_circuit", true)) - if mlx != nil { - mlx.Close(ctx, resp, nil) + // before rate limiting so a non-dispatched request consumes no permit. A + // routed tip leader has the block, so it is not in the future. + if len(tipLeaders) == 0 { + if resp, ok := n.tryShortCircuitFutureBlock(ctx, req, method); ok { + forwardSpan.SetAttributes(attribute.Bool("future_block.short_circuit", true)) + if mlx != nil { + mlx.Close(ctx, resp, nil) + } + return resp, nil } - return resp, nil } // 3) Check if we should handle this method on this network @@ -2301,37 +2313,6 @@ func (n *Network) Forward(ctx context.Context, req *common.NormalizedRequest) (* return requestBlockNumber(ctx, effectiveReq) > 0 && !evm.EmptyResultBeyondConfidence(ctx, effectiveReq) } - // Tip-leader routing: a request pinned to a block above every routed - // upstream's head goes first to the fallbacks that already have it - // (a fallback's WS heads keep its poller current), then to the routed - // list. It takes the per-request escalation, so no other leg or - // attempt escapes; this sweep keeps its escape to the fallbacks it - // has not tried. - leaderRouted := false - if !oneUpstreamOnly && n.cfg.Failover.Enabled() && !failsafeExecutor.HasConsensus() { - if leaders := n.tipLeaderFallbacks(execSpanCtx, effectiveReq, method, upsList); len(leaders) > 0 && - effectiveReq.MarkEscalatedToFallbacks() { - leaderRouted = true - routed := effectiveReq.NextUpstream - maxLoopIterations += len(leaders) - nextUpstream = func() (common.Upstream, error) { - for len(leaders) > 0 { - fb := leaders[0] - leaders = leaders[1:] - if _, loaded := effectiveReq.ConsumedUpstreams.LoadOrStore(fb, true); !loaded { - return fb, nil - } - } - return routed() - } - - telemetry.MetricNetworkTipLeaderRouteTotal.WithLabelValues( - n.projectId, n.Label(), method, - ).Inc() - lg.Debug().Int("leaders", len(leaders)).Msg("routing to fallback upstreams ahead of the routed tip") - } - } - escalationLoop: for { for loopIteration := 0; loopIteration < maxLoopIterations; loopIteration++ { @@ -2530,8 +2511,7 @@ func (n *Network) Forward(ctx context.Context, req *common.NormalizedRequest) (* fallbacks = append(fallbacks, fb) } } - if len(fallbacks) > 0 && (leaderRouted || effectiveReq.MarkEscalatedToFallbacks()) { - leaderRouted = false + if len(fallbacks) > 0 && effectiveReq.MarkEscalatedToFallbacks() { if bestResp != nil { bestResp.Release() bestResp = nil @@ -3666,12 +3646,9 @@ 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, isFallbackTier) -} - -// isFallbackTier reports whether u is tagged tier:fallback. -func isFallbackTier(u common.Upstream) bool { - return u.Config() != nil && u.Config().HasTag(common.TagTierFallback) + return stablePartition(ups, func(u common.Upstream) bool { + return u.Config() != nil && u.Config().HasTag(common.TagTierFallback) + }) } // partitionUpstreamsByLatestBlock moves upstreams whose polled head is known @@ -3732,18 +3709,11 @@ func preferTipLeaderForNearTipGetBlock(ups []common.Upstream, method string, bn return out } -// tipLeaderFallbacks returns the fallback-tier upstreams, outside the routed -// list and allowed by the request's upstream selector, whose head has reached -// the block req is pinned to, when every routed upstream's head is known and -// below it. Nil otherwise. -func (n *Network) tipLeaderFallbacks(ctx context.Context, req *common.NormalizedRequest, method string, routed []common.Upstream) []common.Upstream { - if n.Architecture() != common.ArchitectureEvm { - return nil - } - bn := requestBlockNumber(ctx, req) - if bn <= 0 { - return nil - } +// tipLeaderFallbacks returns the fallback-escape upstreams outside routed, +// allowed by the request's upstream selector, whose polled head has reached bn +// and whose enforced availability bounds admit it, when every routed +// upstream's head is known and below bn. Nil otherwise. +func (n *Network) tipLeaderFallbacks(ctx context.Context, req *common.NormalizedRequest, method string, bn int64, routed []common.Upstream) []common.Upstream { routedIds := make(map[string]struct{}, len(routed)) for _, u := range routed { if lb := upstreamLatestBlock(u); lb <= 0 || lb >= bn { @@ -3751,20 +3721,13 @@ func (n *Network) tipLeaderFallbacks(ctx context.Context, req *common.Normalized } routedIds[u.Id()] = struct{}{} } - return n.fallbacksAtBlock(ctx, req, method, bn, routedIds) -} - -// fallbacksAtBlock returns the fallback-escape upstreams, not in skip and -// allowed by the request's upstream selector, whose head has reached bn and -// whose enforced availability bounds admit it. -func (n *Network) fallbacksAtBlock(ctx context.Context, req *common.NormalizedRequest, method string, bn int64, skip map[string]struct{}) []common.Upstream { selector := "" if d := req.Directives(); d != nil { selector = d.UseUpstream } var out []common.Upstream for _, fb := range n.upstreamsRegistry.GetFallbackEscapeUpstreams(ctx, n.networkId, method) { - if _, ok := skip[fb.Id()]; ok || upstreamLatestBlock(fb) < bn || !n.availabilityAdmits(fb, method, bn) { + if _, ok := routedIds[fb.Id()]; ok || upstreamLatestBlock(fb) < bn || !n.availabilityAdmits(fb, method, bn) { continue } if selector != "" { From df4d7ebdee40d95a8e7bdbaaff94169ae820899c Mon Sep 17 00:00:00 2001 From: Jonny Date: Mon, 5 Oct 2026 22:19:31 +0100 Subject: [PATCH 3/4] revert(network): drop tip-leader routing to fallbacks, keep the hedge fix Fallbacks are only for when the primaries are down. Routing block-pinned reads to a fallback whose head leads sends traffic to it while the primaries are healthy, so the routing, its metric and the future-block short-circuit exception go. The race that motivated it is fixed where it starts instead: eRPC delivering a fallback's head before any primary has the block (next commit). The hedge keeper fix stays: an all-missing ErrUpstreamsExhausted is no longer kept, because a sibling leg may be on an upstream this leg never tried. That already happens in the plain escape path, with no tip leader involved (both TestFailover_Hedge* tests fail without it). Co-Authored-By: Claude Opus 5.5 --- docs/pages/config/projects/networks.mdx | 2 - docs/pages/reference/metrics.mdx | 1 - erpc/networks.go | 87 +--- erpc/networks_failover_escape_test.go | 19 +- erpc/networks_failover_hedge_test.go | 167 ++++++++ erpc/networks_tip_leader_test.go | 544 ------------------------ telemetry/metrics.go | 9 - 7 files changed, 180 insertions(+), 649 deletions(-) create mode 100644 erpc/networks_failover_hedge_test.go delete mode 100644 erpc/networks_tip_leader_test.go diff --git a/docs/pages/config/projects/networks.mdx b/docs/pages/config/projects/networks.mdx index 352da637b..33e967ce3 100644 --- a/docs/pages/config/projects/networks.mdx +++ b/docs/pages/config/projects/networks.mdx @@ -147,8 +147,6 @@ The escape: - fires at most once per request, and never for consensus requests; - is counted in `erpc_network_fallback_escape_total{project,network,category}` — expect zero in steady state. -With `onDefaultsExhausted` on, a request pinned to a block above every routed upstream's known head also routes to the `tier:fallback` upstreams whose head has already reached that block (for example because their `newHeads` subscription announced it first). For any block-pinned request, upstreams whose head has the block are tried before those whose head is below it, whatever their tier. This is ordering, not an escape: the escape above still fires once, to the fallbacks not yet tried, if every upstream fails. Fallbacks are added only when every routed upstream's head is known and below the block, and never for consensus requests; each addition is counted in `erpc_network_tip_leader_route_total{project,network,category}`. - The validation report warns when `onDefaultsExhausted` is enabled but no upstream is tagged `tier:fallback`. ## Agent reference diff --git a/docs/pages/reference/metrics.mdx b/docs/pages/reference/metrics.mdx index 9d0e830b8..b6e97c800 100644 --- a/docs/pages/reference/metrics.mdx +++ b/docs/pages/reference/metrics.mdx @@ -312,7 +312,6 @@ All metric names carry the `erpc_` prefix. Full definitions: 0 && n.cfg.Failover.Enabled() { - if fe := n.getFailsafeExecutor(ctx, req); fe != nil && !fe.HasConsensus() { - if tipLeaders = n.tipLeaderFallbacks(ctx, req, method, bn, upsList); len(tipLeaders) > 0 { - upsList = append(slices.Clone(upsList), tipLeaders...) - telemetry.MetricNetworkTipLeaderRouteTotal.WithLabelValues( - n.projectId, n.Label(), method, - ).Inc() - } + if bn := requestBlockNumber(ctx, req); bn > 0 { + upsList = partitionUpstreamsByLatestBlock(upsList, bn) + upsList = preferTipLeaderForNearTipGetBlock(upsList, method, bn) } } @@ -2065,13 +2056,6 @@ func (n *Network) Forward(ctx context.Context, req *common.NormalizedRequest) (* upsList = tierUpstreamsByGroup(upsList) } - // For a specific block, try upstreams whose poller already has it first, - // whatever their tier. Ordering hints only: every upstream stays eligible. - if bn > 0 { - upsList = partitionUpstreamsByLatestBlock(upsList, bn) - upsList = preferTipLeaderForNearTipGetBlock(upsList, method, bn) - } - // Architecture-specific pruning of the upstream list. Currently only SVM // uses this hook; both filters are gated on ArchitectureSvm so EVM networks // never enter this block. @@ -2169,16 +2153,13 @@ func (n *Network) Forward(ctx context.Context, req *common.NormalizedRequest) (* // Future-block short-circuit: a concrete block number beyond every eligible // upstream's head cannot be served yet — return the truthful null instead of // dispatching + hedging across upstreams that will all return empty. Runs - // before rate limiting so a non-dispatched request consumes no permit. A - // routed tip leader has the block, so it is not in the future. - if len(tipLeaders) == 0 { - if resp, ok := n.tryShortCircuitFutureBlock(ctx, req, method); ok { - forwardSpan.SetAttributes(attribute.Bool("future_block.short_circuit", true)) - if mlx != nil { - mlx.Close(ctx, resp, nil) - } - return resp, nil + // before rate limiting so a non-dispatched request consumes no permit. + if resp, ok := n.tryShortCircuitFutureBlock(ctx, req, method); ok { + forwardSpan.SetAttributes(attribute.Bool("future_block.short_circuit", true)) + if mlx != nil { + mlx.Close(ctx, resp, nil) } + return resp, nil } // 3) Check if we should handle this method on this network @@ -3709,52 +3690,6 @@ func preferTipLeaderForNearTipGetBlock(ups []common.Upstream, method string, bn return out } -// tipLeaderFallbacks returns the fallback-escape upstreams outside routed, -// allowed by the request's upstream selector, whose polled head has reached bn -// and whose enforced availability bounds admit it, when every routed -// upstream's head is known and below bn. Nil otherwise. -func (n *Network) tipLeaderFallbacks(ctx context.Context, req *common.NormalizedRequest, method string, bn int64, routed []common.Upstream) []common.Upstream { - routedIds := make(map[string]struct{}, len(routed)) - for _, u := range routed { - if lb := upstreamLatestBlock(u); lb <= 0 || lb >= bn { - return nil - } - routedIds[u.Id()] = struct{}{} - } - selector := "" - if d := req.Directives(); d != nil { - selector = d.UseUpstream - } - var out []common.Upstream - for _, fb := range n.upstreamsRegistry.GetFallbackEscapeUpstreams(ctx, n.networkId, method) { - if _, ok := routedIds[fb.Id()]; ok || upstreamLatestBlock(fb) < bn || !n.availabilityAdmits(fb, method, bn) { - continue - } - if selector != "" { - if match, err := common.UpstreamMatchesSelector(selector, fb); err != nil || !match { - continue - } - } - out = append(out, fb) - } - return out -} - -// availabilityAdmits reports whether u's block-availability bounds, where -// enforced for method, admit bn. The same bounds checkUpstreamBlockAvailability -// gates on, without its metrics. -func (n *Network) availabilityAdmits(u common.Upstream, method string, bn int64) bool { - if methodHasDedicatedRangeAvailabilityHook(method) || n.blockAvailabilityExplicitlyDisabled(method) { - return true - } - eu, ok := u.(common.EvmUpstream) - if !ok { - return true - } - lo, hi := eu.EvmBlockAvailabilityBounds() - return (lo == math.MinInt64 || bn >= lo) && (hi == math.MaxInt64 || bn <= hi) -} - // upstreamLatestBlock is u's polled head, or 0 when unknown. func upstreamLatestBlock(u common.Upstream) int64 { if eu, ok := u.(common.EvmUpstream); ok { diff --git a/erpc/networks_failover_escape_test.go b/erpc/networks_failover_escape_test.go index 945a59372..97434ac63 100644 --- a/erpc/networks_failover_escape_test.go +++ b/erpc/networks_failover_escape_test.go @@ -276,26 +276,11 @@ func TestFailover_EscapeHatch(t *testing.T) { ctx, cancel := context.WithCancel(context.Background()) defer cancel() - // Primaries at 1002 fail eth_call with a retryable error; the - // fallbacks at 1002 serve it. (Primaries below the block would take - // tip-leader routing instead, see TestFailover_TipLeaderRouting.) + // Primaries at 1000 skip block 1002; the fallbacks at 1002 serve it. network, _, _ := setupFailoverFixture(t, ctx, failoverFixtureOpts{ - primaryLatest: "0x3ea", // 1002 + primaryLatest: "0x3e8", // 1000 fallbackLatest: "0x3ea", // 1002 enableFailover: true, - mocks: func() { - for _, host := range []string{"rpc1.localhost", "rpc2.localhost"} { - host := host - 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). - JSON([]byte(`{"jsonrpc":"2.0","id":1,"error":{"code":-32603,"message":"internal error"}}`)) - } - }, }) counter := telemetry.MetricNetworkFallbackEscapeTotal.WithLabelValues("main", "evm:999", "eth_call") diff --git a/erpc/networks_failover_hedge_test.go b/erpc/networks_failover_hedge_test.go new file mode 100644 index 000000000..cf6d1765f --- /dev/null +++ b/erpc/networks_failover_hedge_test.go @@ -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)) +} diff --git a/erpc/networks_tip_leader_test.go b/erpc/networks_tip_leader_test.go deleted file mode 100644 index d9241f83a..000000000 --- a/erpc/networks_tip_leader_test.go +++ /dev/null @@ -1,544 +0,0 @@ -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" -) - -// unboundedPrimaries drops the availability bound from the primaries, so they -// are called (and answer) for blocks above their polled head, like a local -// node whose poller trails the head a fallback just announced. -func unboundedPrimaries(cfgs []*common.UpstreamConfig) { - for _, cfg := range cfgs { - if !cfg.HasTag(common.TagTierFallback) { - cfg.Evm.BlockAvailability = nil - } - } -} - -// 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) -} - -// missingOnPrimariesSlowOnFallbacks makes the primaries answer eth_call with -// missing data and the fallbacks serve it after delay. -func missingOnPrimariesSlowOnFallbacks(delay time.Duration) func() { - return func() { - for _, host := range []string{"rpc1.localhost", "rpc2.localhost"} { - host := host - 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). - JSON([]byte(`{"jsonrpc":"2.0","id":1,"error":{"code":-32000,"message":"header not found"}}`)) - } - for _, host := range []string{"rpc3.localhost", "rpc4.localhost"} { - host := host - 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(`{"jsonrpc":"2.0","id":1,"result":"0x3333"}`)) - } - } -} - -// hedgeFaster hedges well before a fallback 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 -} - -func TestFailover_TipLeaderRouting(t *testing.T) { - leaderCounter := func() float64 { - return promUtil.ToFloat64(telemetry.MetricNetworkTipLeaderRouteTotal.WithLabelValues("main", "evm:999", "eth_call")) - } - escapeCounter := func() float64 { - return promUtil.ToFloat64(telemetry.MetricNetworkFallbackEscapeTotal.WithLabelValues("main", "evm:999", "eth_call")) - } - - t.Run("RoutesToFallbackThatHasTheBlock", func(t *testing.T) { - defer util.ResetGock() - ctx, cancel := context.WithCancel(context.Background()) - defer cancel() - - // Primaries polled at 1000 would still answer for 1002; the fallbacks - // already have 1002, so they go first. - var primaryHits atomic.Int64 - network, _, _ := setupFailoverFixture(t, ctx, failoverFixtureOpts{ - primaryLatest: "0x3e8", // 1000 - fallbackLatest: "0x3ea", // 1002 - enableFailover: true, - configure: unboundedPrimaries, - mocks: func() { - countEthCalls("rpc1.localhost", &primaryHits) - countEthCalls("rpc2.localhost", &primaryHits) - }, - }) - - leaderBefore, escapeBefore := leaderCounter(), escapeCounter() - result, err := forwardEthCall(t, ctx, network, 1, "0x3ea") - require.NoError(t, err) - assert.Contains(t, []string{"0x3333", "0x4444"}, result, "a fallback that has the block must serve it") - assert.Equal(t, int64(0), primaryHits.Load(), "primaries must not be tried before the leader") - assert.Equal(t, leaderBefore+1, leaderCounter()) - assert.Equal(t, escapeBefore, escapeCounter(), "leader routing is not an escape") - }) - - t.Run("NoLeaderWhenRoutedHasTheBlock", func(t *testing.T) { - defer util.ResetGock() - ctx, cancel := context.WithCancel(context.Background()) - defer cancel() - - network, _, _ := setupFailoverFixture(t, ctx, failoverFixtureOpts{ - primaryLatest: "0x3ea", // 1002 - fallbackLatest: "0x3eb", // 1003 - enableFailover: true, - }) - - leaderBefore := leaderCounter() - for i := 0; i < 10; i++ { - result, err := forwardEthCall(t, ctx, network, i, "0x3ea") - require.NoError(t, err) - assert.Contains(t, []string{"0x1111", "0x2222"}, result, "a primary that has the block must serve it (iter %d)", i) - } - assert.Equal(t, leaderBefore, leaderCounter()) - }) - - t.Run("NoLeaderWhenFallbackIsBehindToo", func(t *testing.T) { - defer util.ResetGock() - ctx, cancel := context.WithCancel(context.Background()) - defer cancel() - - network, _, _ := setupFailoverFixture(t, ctx, failoverFixtureOpts{ - primaryLatest: "0x3e8", // 1000 - fallbackLatest: "0x3e9", // 1001 - enableFailover: true, - configure: unboundedPrimaries, - }) - - leaderBefore := leaderCounter() - result, err := forwardEthCall(t, ctx, network, 1, "0x3ea") - require.NoError(t, err) - assert.Contains(t, []string{"0x1111", "0x2222"}, result) - assert.Equal(t, leaderBefore, leaderCounter()) - }) - - t.Run("NoLeaderWhenFailoverDisabled", func(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: false, - configure: unboundedPrimaries, - }) - - leaderBefore := leaderCounter() - _, _ = forwardEthCall(t, ctx, network, 1, "0x3ea") - assert.Equal(t, leaderBefore, leaderCounter()) - }) - - t.Run("FallsBackToRoutedWhenLeaderFails", func(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: unboundedPrimaries, - mocks: func() { - for _, host := range []string{"rpc3.localhost", "rpc4.localhost"} { - host := host - 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). - JSON([]byte(`{"jsonrpc":"2.0","id":1,"error":{"code":-32603,"message":"internal error"}}`)) - } - }, - }) - - leaderBefore, escapeBefore := leaderCounter(), escapeCounter() - result, err := forwardEthCall(t, ctx, network, 1, "0x3ea") - require.NoError(t, err) - assert.Contains(t, []string{"0x1111", "0x2222"}, result, "routed upstreams must still be tried after the leaders fail") - assert.Equal(t, leaderBefore+1, leaderCounter()) - assert.Equal(t, escapeBefore, escapeCounter(), "the escalation is already spent") - }) - - t.Run("HedgeDoesNotCancelSlowLeader", func(t *testing.T) { - defer util.ResetGock() - ctx, cancel := context.WithCancel(context.Background()) - defer cancel() - - // The hedge leg sweeps the primaries (missing data) while the leader - // leg still waits on the fallback; the leader must win. - network, _, _ := setupFailoverFixture(t, ctx, failoverFixtureOpts{ - primaryLatest: "0x3e8", // 1000 - fallbackLatest: "0x3ea", // 1002 - enableFailover: true, - configure: unboundedPrimaries, - failsafe: hedgeFaster(), - mocks: missingOnPrimariesSlowOnFallbacks(80 * time.Millisecond), - }) - - 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) - } - }) -} - -// 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() - - // Nobody's head reaches 1002, so there is no tip leader: 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: func(cfgs []*common.UpstreamConfig) { - for _, cfg := range cfgs { - cfg.Evm.BlockAvailability = nil - } - }, - failsafe: hedgeFaster(), - mocks: missingOnPrimariesSlowOnFallbacks(80 * time.Millisecond), - }) - - escapeBefore := promUtil.ToFloat64(telemetry.MetricNetworkFallbackEscapeTotal.WithLabelValues("main", "evm:999", "eth_call")) - 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(telemetry.MetricNetworkFallbackEscapeTotal.WithLabelValues("main", "evm:999", "eth_call"))) -} - -// 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)) -} - -// latestHeadMock pins host's polled latest block, ahead of its standard mock. -func latestHeadMock(host, latestHex string) { - gock.New("http://" + host). - Post(""). - Persist(). - Filter(func(r *http.Request) bool { - b := util.SafeReadBody(r) - return r.URL.Host == host && strings.Contains(b, "eth_getBlockByNumber") && strings.Contains(b, `"latest"`) - }). - Reply(200). - JSON([]byte(`{"result":{"number":"` + latestHex + `","timestamp":"0x6702a8f0"}}`)) -} - -func unboundedAll(cfgs []*common.UpstreamConfig) { - for _, cfg := range cfgs { - cfg.Evm.BlockAvailability = nil - } -} - -const missingDataBody = `{"jsonrpc":"2.0","id":1,"error":{"code":-32000,"message":"header not found"}}` - -// 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, no tip leader - 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) - } -} - -// When the tip leader fails, the sweep still escapes to the fallbacks it has -// not tried, even one whose polled head trails the block. -func TestFailover_TipLeaderKeepsEscapeToOtherFallbacks(t *testing.T) { - defer util.ResetGock() - ctx, cancel := context.WithCancel(context.Background()) - defer cancel() - - network, _, _ := setupFailoverFixture(t, ctx, failoverFixtureOpts{ - primaryLatest: "0x3e8", // 1000 - fallbackLatest: "0x3ea", // fallback-1 at 1002 (leader) - enableFailover: true, - configure: unboundedAll, - mocks: func() { - latestHeadMock("rpc4.localhost", "0x3e8") // fallback-2 at 1000 - ethCallMock("rpc1.localhost", 0, missingDataBody) - ethCallMock("rpc2.localhost", 0, missingDataBody) - ethCallMock("rpc3.localhost", 0, `{"jsonrpc":"2.0","id":1,"error":{"code":-32603,"message":"internal error"}}`) - }, - }) - - leader := telemetry.MetricNetworkTipLeaderRouteTotal.WithLabelValues("main", "evm:999", "eth_call") - escape := telemetry.MetricNetworkFallbackEscapeTotal.WithLabelValues("main", "evm:999", "eth_call") - leaderBefore, escapeBefore := promUtil.ToFloat64(leader), promUtil.ToFloat64(escape) - - result, err := forwardEthCall(t, ctx, network, 1, "0x3ea") - require.NoError(t, err) - assert.Equal(t, "0x4444", result, "the untried fallback must still serve the request") - assert.Equal(t, leaderBefore+1, promUtil.ToFloat64(leader)) - assert.Equal(t, escapeBefore+1, promUtil.ToFloat64(escape)) -} - -// A request pinned to an upstream by the use-upstream directive never takes -// the tip-leader route to a fallback outside it. -func TestFailover_TipLeaderRespectsUseUpstream(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: unboundedPrimaries, - }) - - leader := telemetry.MetricNetworkTipLeaderRouteTotal.WithLabelValues("main", "evm:999", "eth_call") - before := promUtil.ToFloat64(leader) - - req := ethCallRequest(1, "0x3ea") - req.SetDirectives(&common.RequestDirectives{UseUpstream: "primary-1"}) - req.SetNetwork(network) - resp, err := network.Forward(ctx, req) - require.NoError(t, err) - require.NotNil(t, resp) - defer resp.Release() - jrr, err := resp.JsonRpcResponse() - require.NoError(t, err) - assert.Equal(t, "0x1111", strings.Trim(jrr.GetResultString(), `"`)) - assert.Equal(t, before, promUtil.ToFloat64(leader)) -} - -// With served-tip on, a numbered eth_getBlockByNumber above every eligible -// head is short-circuited to null, unless a reachable fallback already has -// the block. -func TestFailover_FutureBlockShortCircuitSparesFallbackThatHasIt(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, - network: func(cfg *common.NetworkConfig) { - cfg.Evm.ServedTip = &common.EvmServedTipConfig{EnabledFor: []string{"latest"}} - }, - mocks: func() { - for _, host := range []string{"rpc3.localhost", "rpc4.localhost"} { - host := host - gock.New("http://" + host). - Post(""). - Persist(). - Filter(func(r *http.Request) bool { - b := util.SafeReadBody(r) - return r.URL.Host == host && strings.Contains(b, "eth_getBlockByNumber") && strings.Contains(b, `"0x3ea"`) - }). - Reply(200). - JSON([]byte(`{"jsonrpc":"2.0","id":1,"result":{"number":"0x3ea","hash":"0xfb","timestamp":"0x6702a8f2"}}`)) - } - }, - }) - - req := common.NewNormalizedRequest([]byte(`{"jsonrpc":"2.0","id":1,"method":"eth_getBlockByNumber","params":["0x3ea",false]}`)) - req.SetNetwork(network) - resp, err := network.Forward(ctx, req) - require.NoError(t, err) - require.NotNil(t, resp) - defer resp.Release() - jrr, err := resp.JsonRpcResponse() - require.NoError(t, err) - assert.Contains(t, jrr.GetResultString(), `"0x3ea"`, "the fallback's block must be returned, not a synthetic null") -} - -// servedTipLatest turns served-tip on for the latest axis. -func servedTipLatest(cfg *common.NetworkConfig) { - cfg.Evm.ServedTip = &common.EvmServedTipConfig{EnabledFor: []string{"latest"}} -} - -func getBlockRequest(blockHex string) *common.NormalizedRequest { - return common.NewNormalizedRequest([]byte(`{"jsonrpc":"2.0","id":1,"method":"eth_getBlockByNumber","params":["` + blockHex + `",false]}`)) -} - -// Consensus never reaches a fallback, so a fallback that has the block must -// not stop the future-block short-circuit for it. -func TestFailover_FutureBlockShortCircuitStillAppliesToConsensus(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, - network: servedTipLatest, - failsafe: []*common.FailsafeConfig{{ - MatchMethod: "*", - Consensus: &common.ConsensusPolicyConfig{MaxParticipants: 2, AgreementThreshold: 2}, - }}, - }) - - req := getBlockRequest("0x3ea") - req.SetNetwork(network) - resp, err := network.Forward(ctx, req) - require.NoError(t, err, "consensus must still get the truthful null") - require.NotNil(t, resp) - defer resp.Release() - assert.True(t, resp.IsResultEmptyish(), "expected the short-circuit null") -} - -// A fallback whose configured availability excludes the block is neither a -// tip leader nor a reason to skip the future-block short-circuit. -func TestFailover_FallbackAvailabilityBoundsExcludeLeader(t *testing.T) { - leader := telemetry.MetricNetworkTipLeaderRouteTotal.WithLabelValues("main", "evm:999", "eth_call") - capFallbacks := func(cfgs []*common.UpstreamConfig) { - unboundedPrimaries(cfgs) - for _, cfg := range cfgs { - if cfg.HasTag(common.TagTierFallback) { - cfg.Evm.BlockAvailability = &common.EvmBlockAvailabilityConfig{ - Upper: &common.EvmAvailabilityBoundConfig{ExactBlock: i64(1000)}, - } - } - } - } - - t.Run("NotALeader", func(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, but capped at 1000 - enableFailover: true, - configure: capFallbacks, - }) - - before := promUtil.ToFloat64(leader) - result, err := forwardEthCall(t, ctx, network, 1, "0x3ea") - require.NoError(t, err) - assert.Contains(t, []string{"0x1111", "0x2222"}, result) - assert.Equal(t, before, promUtil.ToFloat64(leader)) - }) - - t.Run("ShortCircuitStillApplies", func(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, but capped at 1000 - enableFailover: true, - network: servedTipLatest, - configure: func(cfgs []*common.UpstreamConfig) { - capFallbacks(cfgs) - for _, cfg := range cfgs { - if !cfg.HasTag(common.TagTierFallback) { - cfg.Evm.BlockAvailability = headBoundedAvailability() - } - } - }, - }) - - req := getBlockRequest("0x3ea") - req.SetNetwork(network) - resp, err := network.Forward(ctx, req) - require.NoError(t, err) - require.NotNil(t, resp) - defer resp.Release() - assert.True(t, resp.IsResultEmptyish(), "expected the short-circuit null") - }) -} diff --git a/telemetry/metrics.go b/telemetry/metrics.go index 73fc3968e..033c525e5 100644 --- a/telemetry/metrics.go +++ b/telemetry/metrics.go @@ -759,15 +759,6 @@ var ( Help: "Total number of times the per-request fallback escape hatch fired because the primary upstream set was exhausted with retryable errors.", }, []string{"project", "network", "category"}) - // MetricNetworkTipLeaderRouteTotal counts block-pinned requests sent to a - // fallback-tier upstream first because it already had the block while - // every routed upstream's head was below it. - MetricNetworkTipLeaderRouteTotal = DefineCounter(prometheus.CounterOpts{ - Namespace: "erpc", - Name: "network_tip_leader_route_total", - Help: "Total number of block-pinned requests routed to a fallback-tier upstream first because it had the block and no routed upstream did.", - }, []string{"project", "network", "category"}) - // MetricCacheExecutorAttempt counts every attempt a cache-connector // failsafe executor governs, keyed by the executor's identity (the // matchMethod / matchFinality it was configured with) and the outcome From a228f52b8efc2b2b08508ce98aacb9596d450a24 Mon Sep 17 00:00:00 2001 From: Jonny Date: Tue, 6 Oct 2026 04:37:32 +0100 Subject: [PATCH 4/4] fix(ws): hold back fallback heads while a primary streams heads newHeads is subscribed on every WebSocket upstream and each head goes out from whichever source announces it first. With failover on, a fallback-tier upstream whose WebSocket is faster told clients about a block no primary had yet; their follow-up reads then missed on the primaries, or escaped to the fallback while the primaries were healthy. A fallback's head now reaches clients, and feeds the delivered-head floor, only while no upstream outside that tier which the selection policy routes to has a live newHeads subscription of its own. "Down" is the policy's verdict, as for reads, so custom policies that keep fallbacks in the ordered list behave the same. A network whose only WebSocket upstreams are fallbacks, or whose primaries' subscriptions are dead, keeps receiving fallback heads. A held-back head still updates that fallback's state poller. NetworkHandle.SuggestLatestBlock now reports whether the head may be delivered; the indexer drops it before dedup otherwise. Co-Authored-By: Claude Opus 5.5 --- docs/pages/operation/websocket.mdx | 4 +- erpc/networks.go | 9 +- erpc/networks_ws_tip_test.go | 99 +++++++++++++++++++ erpc/subscription_manager.go | 65 +++++++++--- indexer/adapters/wsupstream/adapter.go | 8 ++ .../wsupstream/adapter_reconnect_test.go | 6 +- indexer/indexer.go | 4 +- indexer/indexer_test.go | 34 ++++++- indexer/ingress.go | 9 +- indexer/integration_test.go | 3 +- 10 files changed, 210 insertions(+), 31 deletions(-) diff --git a/docs/pages/operation/websocket.mdx b/docs/pages/operation/websocket.mdx index 3408a61b4..6fd0aa9a4 100644 --- a/docs/pages/operation/websocket.mdx +++ b/docs/pages/operation/websocket.mdx @@ -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. @@ -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. diff --git a/erpc/networks.go b/erpc/networks.go index be977e951..d1ad6ebd6 100644 --- a/erpc/networks.go +++ b/erpc/networks.go @@ -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 diff --git a/erpc/networks_ws_tip_test.go b/erpc/networks_ws_tip_test.go index d86d83723..bd3c69cf1 100644 --- a/erpc/networks_ws_tip_test.go +++ b/erpc/networks_ws_tip_test.go @@ -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" @@ -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) { diff --git a/erpc/subscription_manager.go b/erpc/subscription_manager.go index b4e0ac8d2..251e185b8 100644 --- a/erpc/subscription_manager.go +++ b/erpc/subscription_manager.go @@ -259,14 +259,22 @@ func (sm *SubscriptionManager) bootstrapNetwork(ctx context.Context, nw *Network return common.NewErrNoWsUpstreamAvailable(networkID) } - sm.idx.RegisterNetwork(&networkHandle{nw: nw}) - sm.idx.RegisterNetworkSelector(networkID, &subIngressSelector{nw: nw}) var adapterOpts wsupstream.Options if cfg := nw.cfg; cfg != nil && cfg.Evm != nil && cfg.Evm.StripSubscribeFromBlockZero != nil { adapterOpts.StripSubscribeFromBlockZero = *cfg.Evm.StripSubscribeFromBlockZero } + adapters := make(map[string]*wsupstream.Adapter, len(wsUpstreams)) + heads := make(map[string]headSource, len(wsUpstreams)) for _, up := range wsUpstreams { - adapter := wsupstream.New(up, networkID, sm.logger, adapterOpts) + if adapter := wsupstream.New(up, networkID, sm.logger, adapterOpts); adapter != nil { + adapters[up.Id()] = adapter + heads[up.Id()] = adapter + } + } + sm.idx.RegisterNetwork(&networkHandle{nw: nw, heads: heads}) + sm.idx.RegisterNetworkSelector(networkID, &subIngressSelector{nw: nw}) + for _, up := range wsUpstreams { + adapter := adapters[up.Id()] if adapter == nil { continue } @@ -416,38 +424,65 @@ func (sm *SubscriptionManager) recordFailureMetrics( // networkHandle adapts *Network to indexer.NetworkHandle. type networkHandle struct { nw *Network + // heads are the network's WebSocket ingresses by upstream id; fixed once + // the network is registered. + heads map[string]headSource +} + +// headSource is an ingress that streams newHeads. +type headSource interface { + // HeadsLive reports whether its newHeads subscription is live. + HeadsLive() bool } func (h *networkHandle) Id() string { return h.nw.networkId } // SuggestLatestBlock passes a head from the ingress "ws:" to -// that upstream's state poller. Once the poller has accepted it (a major -// jump is verified asynchronously first), a head from a tip candidate also -// advances the network's delivered-head floor before clients see it; a -// fallback-tier or policy-excluded upstream cannot lift "latest". -func (h *networkHandle) SuggestLatestBlock(sourceId string, blockNumber int64) { +// that upstream's state poller and reports whether clients may receive it +// (see deliversHeadsFrom). Once the poller has accepted a deliverable head (a +// major jump is verified asynchronously first), it also advances the +// network's delivered-head floor before clients see it. +func (h *networkHandle) SuggestLatestBlock(sourceId string, blockNumber int64) bool { upstreamID, ok := strings.CutPrefix(sourceId, "ws:") if !ok { - return + return true } ctx := context.Background() for _, u := range h.nw.upstreamsRegistry.GetNetworkUpstreams(ctx, h.nw.networkId) { if u.Id() != upstreamID { continue } + deliver := h.deliversHeadsFrom(ctx, u) poller := u.EvmStatePoller() if poller == nil || poller.IsObjectNull() { - return + return deliver } poller.SuggestLatestBlock(blockNumber) - // Every source's heads reach clients (whichever delivers first), so - // every source feeds the floor — once its poller accepted the head, - // which excludes jumps still pending chain-id verification. - if poller.LatestBlock() >= blockNumber { + if deliver && poller.LatestBlock() >= blockNumber { h.nw.NoteObservedLatestBlock(h.nw.appCtx, blockNumber) } - return + return deliver + } + return true +} + +// deliversHeadsFrom reports whether clients may receive u's heads. With +// failover on, a fallback-tier upstream's heads are held back while some +// upstream outside that tier, which the selection policy routes to, has a +// live newHeads subscription of its own: clients are not told of a block +// only a fallback has while the primaries are up. Whether they are up is the +// policy's verdict, as for reads. Without such a primary (none eligible, or +// none streaming heads) the fallbacks' heads are delivered. +func (h *networkHandle) deliversHeadsFrom(ctx context.Context, u common.Upstream) bool { + if h.nw.cfg == nil || !h.nw.cfg.Failover.Enabled() || !isFallbackTier(u) { + return true + } + for _, c := range h.nw.tipCandidateUpstreams(ctx, "*") { + if s := h.heads[c.Id()]; s != nil && !isFallbackTier(c) && s.HeadsLive() { + return false + } } + return true } var ( diff --git a/indexer/adapters/wsupstream/adapter.go b/indexer/adapters/wsupstream/adapter.go index d0814ea9b..0f3caa545 100644 --- a/indexer/adapters/wsupstream/adapter.go +++ b/indexer/adapters/wsupstream/adapter.go @@ -116,6 +116,14 @@ func New(up *upstream.Upstream, networkID string, logger *zerolog.Logger, opts O func (a *Adapter) Name() string { return "ws:" + a.upstreamID } +// HeadsLive reports whether the newHeads subscription is live on the +// current connection. +func (a *Adapter) HeadsLive() bool { + a.subsMu.Lock() + defer a.subsMu.Unlock() + return a.heads.id != "" +} + // Start registers the connection hooks and subscribes newHeads in the // background, detached from ctx. func (a *Adapter) Start(_ context.Context, nw indexer.NetworkHandle, sink indexer.Sink) error { diff --git a/indexer/adapters/wsupstream/adapter_reconnect_test.go b/indexer/adapters/wsupstream/adapter_reconnect_test.go index 672e152c0..ec46a84c7 100644 --- a/indexer/adapters/wsupstream/adapter_reconnect_test.go +++ b/indexer/adapters/wsupstream/adapter_reconnect_test.go @@ -26,8 +26,8 @@ import ( type fakeNetworkHandle struct{} -func (fakeNetworkHandle) Id() string { return "evm:123" } -func (fakeNetworkHandle) SuggestLatestBlock(string, int64) {} +func (fakeNetworkHandle) Id() string { return "evm:123" } +func (fakeNetworkHandle) SuggestLatestBlock(string, int64) bool { return true } type fakeSink struct { events chan indexer.StreamEvent @@ -235,10 +235,12 @@ func TestAdapterResubscribesWithRetryAfterReconnect(t *testing.T) { return send(ctx, nq, bypass) } + require.False(t, a.HeadsLive(), "no heads before the subscribe succeeds") require.NoError(t, a.Start(context.Background(), fakeNetworkHandle{}, sink)) require.Eventually(t, subscribedHeads(a), 3*time.Second, 10*time.Millisecond, "adapter never subscribed newHeads despite the failures stopping after %d", failuresPerEpoch) + require.True(t, a.HeadsLive()) require.GreaterOrEqual(t, forwardCalls.Load(), int64(failuresPerEpoch+1), "expected the subscribe to be retried through failures") diff --git a/indexer/indexer.go b/indexer/indexer.go index e382a5d88..11a592f3e 100644 --- a/indexer/indexer.go +++ b/indexer/indexer.go @@ -303,7 +303,9 @@ func (i *Indexer) Ingest(ev StreamEvent) { // Before dedup, so every source's head counts even if another source // already delivered it. if ev.Kind == KindNewHead && !ev.Block.Zero() && ev.SourceId != "" { - ns.handle.SuggestLatestBlock(ev.SourceId, ev.Block.Number) + if !ns.handle.SuggestLatestBlock(ev.SourceId, ev.Block.Number) { + return + } } // The upstream's removed flag is trusted as-is. diff --git a/indexer/indexer_test.go b/indexer/indexer_test.go index fcbe50949..3d0866c88 100644 --- a/indexer/indexer_test.go +++ b/indexer/indexer_test.go @@ -18,7 +18,8 @@ type fakeNetwork struct { id string mu sync.Mutex - suggestedBy map[string][]int64 // sourceId -> block nums seen + suggestedBy map[string][]int64 // sourceId -> block nums seen + held map[string]struct{} // sourceIds whose heads are not delivered } func newFakeNetwork(id string) *fakeNetwork { @@ -29,10 +30,12 @@ func newFakeNetwork(id string) *fakeNetwork { } func (n *fakeNetwork) Id() string { return n.id } -func (n *fakeNetwork) SuggestLatestBlock(sourceId string, block int64) { +func (n *fakeNetwork) SuggestLatestBlock(sourceId string, block int64) bool { n.mu.Lock() + defer n.mu.Unlock() n.suggestedBy[sourceId] = append(n.suggestedBy[sourceId], block) - n.mu.Unlock() + _, held := n.held[sourceId] + return !held } type fakeEgress struct { @@ -271,6 +274,31 @@ func TestIndexer_NewHead_StalerDroppedKeepsStatePollerFed(t *testing.T) { } } +// A head the network does not deliver is still observed, and does not +// advance the dedup marker past the heads the other sources deliver. +func TestIndexer_NewHead_HeldSourceObservedNotDelivered(t *testing.T) { + idx := newIndexer(t) + nw := newFakeNetwork("evm:1") + nw.held = map[string]struct{}{"ws:fb": {}} + idx.RegisterNetwork(nw) + eg := &fakeEgress{name: "eg1", acceptAllHeads: true} + idx.Attach(eg) + + idx.Ingest(StreamEvent{Kind: KindNewHead, NetworkId: "evm:1", SourceId: "ws:fb", Block: BlockRef{Number: 101, Hash: "0xBBB"}}) + if got := eg.count(); got != 0 { + t.Fatalf("held head delivered: got %d", got) + } + idx.Ingest(StreamEvent{Kind: KindNewHead, NetworkId: "evm:1", SourceId: "ws:up1", Block: BlockRef{Number: 101, Hash: "0xBBB"}}) + if got := eg.count(); got != 1 { + t.Fatalf("the delivering source's head must go out: want 1, got %d", got) + } + nw.mu.Lock() + defer nw.mu.Unlock() + if len(nw.suggestedBy["ws:fb"]) != 1 { + t.Fatalf("held source must still be observed, got %v", nw.suggestedBy) + } +} + func TestIndexer_Log_RefcountFanOutAndTeardown(t *testing.T) { idx := newIndexer(t) nw := newFakeNetwork("evm:1") diff --git a/indexer/ingress.go b/indexer/ingress.go index 3df2d4086..26338e672 100644 --- a/indexer/ingress.go +++ b/indexer/ingress.go @@ -11,10 +11,11 @@ type Sink interface { // NetworkHandle is the part of a network an ingress may touch. type NetworkHandle interface { Id() string - // SuggestLatestBlock reports a head observed by source. It is called - // before dedup and fan-out, so every source's observation counts and - // the network tip moves before clients see the head. - SuggestLatestBlock(sourceId string, blockNumber int64) + // SuggestLatestBlock reports a head observed by source and whether + // clients may receive it. It is called before dedup and fan-out, so + // every source's observation counts and the network tip moves before + // clients see the head; a head it does not deliver is dropped. + SuggestLatestBlock(sourceId string, blockNumber int64) (deliver bool) } // EventIngress turns a transport-specific subscription into StreamEvents diff --git a/indexer/integration_test.go b/indexer/integration_test.go index 976144854..ebecd78a8 100644 --- a/indexer/integration_test.go +++ b/indexer/integration_test.go @@ -66,13 +66,14 @@ type stubNetwork struct { } func (s *stubNetwork) Id() string { return s.id } -func (s *stubNetwork) SuggestLatestBlock(sourceID string, block int64) { +func (s *stubNetwork) SuggestLatestBlock(sourceID string, block int64) bool { s.mu.Lock() if s.suggestions == nil { s.suggestions = make(map[string][]int64) } s.suggestions[sourceID] = append(s.suggestions[sourceID], block) s.mu.Unlock() + return true } type recordingEgress struct {