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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 6 additions & 0 deletions common/request.go
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down
2 changes: 2 additions & 0 deletions docs/pages/config/projects/networks.mdx
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
1 change: 1 addition & 0 deletions docs/pages/reference/metrics.mdx
Original file line number Diff line number Diff line change
Expand Up @@ -312,6 +312,7 @@ All metric names carry the `erpc_` prefix. Full definitions: <SourceLink file="t
| `erpc_network_successful_request_total` | counter | project, network, vendor, upstream, category, attempt, finality, emptyish, user, agent_name | Request succeeded. `emptyish` ∈ `"true"`/`"false"` — true means the response was empty-ish (null, empty array). |
| `erpc_network_multiplexed_request_total` | counter | project, network, category, finality, user, agent_name | Request de-duplicated into an identical in-flight request. |
| `erpc_network_fallback_escape_total` | counter | project, network, category | The per-request fallback escape (`failover.onDefaultsExhausted`) retried a request on `tier:fallback` upstreams. Expected zero in steady state. |
| `erpc_network_tip_leader_route_total` | counter | project, network, category | A block-pinned request was sent to `tier:fallback` upstreams first because they had the block and no routed upstream did (`failover.onDefaultsExhausted`). |
| `erpc_upstream_websocket_connected` | gauge | project, vendor, network, upstream | `1` while the upstream's WebSocket connection is established, `0` while it reconnects. |
| `erpc_websocket_subscription_notifications_dropped_total` | counter | project, network, kind | Client subscription notifications dropped because the subscription's buffer (`server.webSocket.subscriptionBufferSize`) was full. For `kind="logs"` the connection is also closed with `1013`. |
| `erpc_network_static_response_served_total` | counter | project, network, category | Served from a configured static response; no upstream touched. |
Expand Down
5 changes: 4 additions & 1 deletion erpc/network_executor.go
Original file line number Diff line number Diff line change
Expand Up @@ -668,8 +668,11 @@ func (e *networkExecutor) runHedge(
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 true
return !req.EscalatedToFallbacks()
}
}
// Underlying-retryable wrapped errors (e.g. ErrUpstreamsExhausted
Expand Down
110 changes: 106 additions & 4 deletions erpc/networks.go
Original file line number Diff line number Diff line change
Expand Up @@ -915,6 +915,13 @@ 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
Expand Down Expand Up @@ -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++ {
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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 {
Expand Down
32 changes: 29 additions & 3 deletions erpc/networks_failover_escape_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -169,6 +169,9 @@ func buildFailoverNetwork(
networkConfig.Failover = &common.FailoverConfig{OnDefaultsExhausted: util.BoolPtr(true)}
}
networkConfig.Failsafe = opts.failsafe
if opts.network != nil {
opts.network(networkConfig)
}

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

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

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

upsList := upr.GetNetworkUpstreams(ctx, util.EvmNetworkId(999))
require.Len(t, upsList, 4)
Expand Down Expand Up @@ -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")
Expand Down
Loading
Loading