Skip to content
Open
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
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 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
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
15 changes: 0 additions & 15 deletions erpc/network_executor.go
Original file line number Diff line number Diff line change
Expand Up @@ -656,21 +656,6 @@ func (e *networkExecutor) runHedge(
if uxe.Upstreams() == nil || len(uxe.Upstreams()) == 0 {
return false
}
// When every upstream returned -32004/missing-data, no sibling
// hedge leg can do better — they will all find the same upstreams
// consumed. Keep this result so the retry layer (shouldRetryWithReason)
// sees ErrUpstreamsExhausted directly and applies the 500ms delay.
causes := uxe.Errors()
allMissing := len(causes) > 0
for _, c := range causes {
if !common.HasErrorCode(c, common.ErrCodeEndpointMissingData) {
allMissing = false
break
}
}
if allMissing {
return true
}
}
// Underlying-retryable wrapped errors (e.g. ErrUpstreamsExhausted
// wrapping a 5xx) should continue racing for a healthier sibling.
Expand Down
87 changes: 76 additions & 11 deletions erpc/networks.go
Original file line number Diff line number Diff line change
Expand Up @@ -2041,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()
}
}
}

Expand All @@ -2056,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.
Expand Down Expand Up @@ -2153,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
Expand Down Expand Up @@ -3690,6 +3709,52 @@ 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 {
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