From 8e98cb28b67cc881d67456371225dc0696f695f1 Mon Sep 17 00:00:00 2001 From: shpookas Date: Wed, 30 Sep 2026 09:04:53 +0200 Subject: [PATCH] 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