Skip to content

fix(evm): skip lagging upstreams for eth_call(latest) and mirror head lag to */All - #25

Open
1marcghannam wants to merge 45 commits into
feat/websocket-supportfrom
fix/eth-call-latest-skip-lagging-upstreams
Open

1marcghannam wants to merge 45 commits into
feat/websocket-supportfrom
fix/eth-call-latest-skip-lagging-upstreams

Conversation

@1marcghannam

@1marcghannam 1marcghannam commented Sep 15, 2026 •

Copy link
Copy Markdown
Member

Summary

  • Incident: Priority Pool totalQueued via eRPC returned 16,355 instead of ~692 because a stalled WS peer (internal-eth-mainnet-reth-ws-0) still answered eth_call("latest"). Tip-pin image only covers getBlock.
  • Root cause: eth_call("latest") resolves to bn=0 and fail-opens availability gating; {*,All} BlockHeadLag was not mirrored for peer upstreams so network-scope lag predicates can no-op (upstream #1019).
  • Fix (on tip-pin lineage):
    1. Hard gate: skip upstreams lagging TipHW by >16 for latest-state methods + latest/pending/empty tag.
    2. Tip partition: prefer near-tip peers first for those reads (failover still allowed).
    3. Tracker fix(health): lag predicates silently no-op for CB-open upstreams under evalScope:network erpc/erpc#1019 port: mirror BlockHeadLag into {ups,"*",All}.

Base: fix/pin-near-tip-getblock-to-evm-leader (LinkPool tip-pin / WS lineage — not upstream main).

Test plan

  • TestNetwork_EthCallLatest_SkipsLaggingUpstream
  • TestNetwork_EthCallConcreteBlock_AllowsLaggingArchive
  • TestWildcardLagMirroredForPeerUpstreams
  • Near-tip getBlock pin still green

Rollout

  1. Build/publish image from this tip-pin + eth_call lineage.
  2. Bump internal-erpc image tag.
  3. Verify Priority Pool totalQueued via eRPC ≈ 692; no stalled reth-0 while behind TipHW.

jleeh and others added 30 commits June 1, 2026 12:00
Adds a WebSocket transport for both inbound clients and outbound
upstreams, plus a transport-agnostic event indexer that fans head and
log notifications from any subscribed ingress out to any number of
subscribed clients.

Inbound (client-facing):
- Server-side WS listener served on the same HTTP port. Each client
  connection is a long-lived JSON-RPC session that supports both
  regular RPCs (eth_call, eth_blockNumber, …) and subscriptions
  (eth_subscribe newHeads / logs / eth_unsubscribe). Directives
  (UseUpstream, headers etc.) are honoured via query params and
  upgrade headers.
- Subscriptions are dedup'd in the indexer, so N clients subscribing
  to the same logs filter result in a single upstream subscription.
- Failed subscribes do not close the client connection; only the
  failing eth_subscribe call returns an error.

Outbound (upstream-facing):
- WsJsonRpcClient parallels HttpJsonRpcClient and is wired into the
  same upstream selection / failsafe / scoring / rate-limit / metrics
  pipeline. eth_subscribe is dispatched here; regular RPCs route to
  whichever client (HTTP or WS) the upstream advertises.
- Outbound JSON-RPC ids on the WS wire are rewritten to a unique
  per-client atomic counter so concurrent requests sharing a caller-
  supplied id (small ints are common) do not cross-wire each other in
  the pending response map. The caller's original id is restored on
  the response before returning.
- A NodeGroup-aware "share block-head state across siblings" hook is
  not included in this commit; can be revisited as a follow-up.

Indexer:
- Indexer fans events from registered ingresses out to subscribed
  egresses. Two ingress adapters ship: wsupstream (wraps an upstream
  WS client) and nullingress (no-op for tests).
- newHead dedup is single-CAS on a packed (num, hash) atomic.Pointer
  so concurrent ingest of the same head doesn't double-deliver.
- Reorg handling rolls subscriptions back to the canonical chain on
  reorg-style notifications.
- A `stripSubscribeFromBlockZero` network-level flag drops the
  meaningless `fromBlock: "0x0"` field from `eth_subscribe("logs")`
  filters before forwarding, so live-stream subscribes succeed on
  pruned backends. Non-zero fromBlocks pass through.

Documentation, schema generation, and TS bindings are updated
accordingly.
…efresh

Two related fixes to the shared-state subscriber path.

1) keep-latest pub/sub delivery. notifySubscribers used a try-send-or-
skip pattern against a cap=1 subscriber channel. Under any burst
(messageLoop + pollingLoop racing, or two instances publishing close
together faster than the consumer drains), the newer value was
silently dropped. For monotonic counters (latest / finalized block),
the fresh value is the only one correctness depends on — dropping it
reopens the propagation window the shared counter exists to close
and produces cross-instance regressions that repeatable-read clients
treat as data corruption.

Replaced with sendKeepLatest: on a full buffer, drain the stale entry
and retry the send. The counter is monotonic so displacing older by
newer is always right.

At the same time, drop the mutex-around-channel-send anti-pattern the
old subscriber struct carried. A `done` channel closed once via
sync.Once signals subscriber shutdown; writers select on
ch | done | default and never need a lock. Receivers that want the
shutdown signal select on done — sc.ch is no longer closed by any
path. Cleanup and stop-all-subscribers both call sc.close(),
idempotent.

2) honor caller ctx deadline for refresh fn timeout. The refresh
function was running with the registry's app-level context, which
ignored per-call deadlines. Plumbing the caller's context through
makes refresh respect timeouts the consumer sets, instead of running
unbounded.

Regression and race-detector coverage in
redis_pubsub_manager_test.go and shared_state_variable_test.go.
A heap/goroutine profile from a long-running deployment showed configured
upstream counts expanding by 5-10x at runtime: the per-upstream client
instances (HTTP and WS), EvmStatePoller polling goroutines, and
sharedStateRegistry counter-sync goroutines all multiplied past their
configured set, with heap retention concentrated in upstream.NewUpstream
and its policy/health-tracker builders.

Four interlocking bugs were causing transient and ostensibly-deduplicated
upstream constructions to actually leak:

1. common.UniqueUpstreamKey hashed up.NetworkId() (which returns "n/a"
   before registration and the real id after) and iterated headers in
   non-deterministic Go map order. The same upstream produced different
   keys at different points in its lifecycle, defeating the per-upstream
   client cache and shared-state counter dedup. Drop NetworkId from the
   key (cfg.Endpoint already disambiguates) and sort header keys before
   hashing. Regression test in common/upstream_test.go.

2. Upstream.Bootstrap unconditionally reassigned u.evmStatePoller =
   evm.NewEvmStatePoller(...) and called Bootstrap on it. The new poller
   spawns a goroutine that listens only on appCtx, so any re-entry of
   Bootstrap orphans the prior poller's goroutine for the lifetime of
   the process. Guard creation with evmStatePollerMu + nil-check so
   repeated Bootstrap calls are idempotent for the poller.

3. ClientRegistry.CreateClient declared `var once sync.Once` locally,
   so every call ran the body. Concurrent callers that both missed the
   `clients` cache could each spawn an HTTP/WS client (with shutdown
   waiter, ping/read loops); only one won the manager.clients.Store —
   the losers' goroutines leaked. Replace with a sync.Map[key]
   *clientCreation that shares both the once and the build result, so
   concurrent callers all return the winning client.

4. config_analyzer.validateUpstreamEndpoints and GenerateValidationReport
   passed the caller's ctx (effectively appCtx) into upstream.NewUpstream
   for transient validation upstreams. Their client goroutines listen on
   that ctx and so outlived validation and accumulated forever. Wrap
   each validation pass in a context.WithCancel and pass the child ctx
   to NewUpstream; the deferred cancel terminates the transient client
   goroutines on validation completion.

After these fixes the per-upstream goroutine multiplier collapses back
to ~1 across all four populations (client, shutdown waiter, poller,
counter-sync), and heap retention from leaked NewUpstream instances
goes with it.

Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
The wrapper previously closed only the gzip.Reader and returned it to the
pool; the underlying io.ReadCloser source was never closed. Single-request
HTTP client callers in clients/http_json_rpc_client.go attached resp.Body
via WrapGzipReader and relied on the wrapper to close it later, but that
close never propagated to the http transport.

Concretely, the docstring on the call site —
  "DO NOT close resp.Body here - it will be closed by NormalizedResponse
   after reading"
was only true for non-gzipped responses (where bodyReader = resp.Body
directly). For gzipped responses the resp.Body was never closed at all
once the gzip wrapper was the head of the chain. Over time that leaks
http2ClientConn.readLoop goroutines as net/http keeps opening new
connections instead of reusing pinned ones.

Have the wrapper take an optional source Closer and close it from Close
after the gzip.Reader is returned to the pool. Update the single-request
caller to pass resp.Body. The batch-path readResponseBody is unaffected
because it already manages resp.Body lifetime via an explicit defer.

Unit tests cover: source is closed exactly once, Close is idempotent under
contention, nil source is accepted, source-close errors propagate.
…uard

A prior cleanup pass dropped the lastReturned*Block atomic.Int64 fields
on Network and the matching CAS-bump + WARN logic at the tail of
evmHighestBlockNumber. Without it, a downstream client polling
EvmHighestFinalizedBlockNumber through a load balancer can observe a
returned value lower than a previously returned value (e.g. when the
aggregator transiently picks up a lagging upstream after a restart, or
when the cross-cluster shared counter is briefly stale relative to the
local primaryMax). Strict downstream consumers treat that as a
data-integrity violation, even though the underlying chain is fine.

Restore the per-tag atomic high-water mark and *clamp* regressions
back to that high-water mark instead of returning the regressed value
with only a warning. Clamping is safe because we never invent a number
— the floor is always something this Network has already returned
upstream. On forward progress we CAS-bump the high-water mark; on
regression we log the inputs (localMax, primaryMax, fallbackMax,
anyPrimaryUp, sharedVal, sharedNil) so future regressions are
diagnosable, and return the previously-returned value.
…escape

Reproducing a production node-down incident confirmed that eRPC fails to
fall over to configured fallback upstreams. Two paths silently bypassed
the health tracker and a third gap made true HA impossible:

1. checkUpstreamBlockAvailability rejections (handleBlockSkip) updated
   Prometheus counters but never called RecordUpstreamRequest /
   RecordUpstreamFailure. An upstream gate-rejected on every request
   showed errorRate=0 to the JS selectionPolicy.

2. PollLatestBlockNumber / PollFinalizedBlockNumber fetch failures
   logged at warn and returned without recording. The CB counted them
   (which is why it cycled half-open <-> open during the outage), but
   subsequent CB-open responses bypass tryForward and the tracker stays
   starved of signal.

3. Even after Fixes (1) and (2), the selectionPolicy's reaction is
   bounded by evalInterval (1m in typical configs) + scoreRefreshInterval
   (10s default). Worst case ~70s of ErrUpstreamsExhausted to clients
   during the transition. Real HA requires sub-request escalation.

This commit lands all three architectural fixes:

* Fix 2 - networks.go:handleBlockSkip records (request, failure) when
  the block-availability rejection is retryable.

* Fix 3 - evm_state_poller.go records (request, failure) when fetchBlock
  returns a CB-open short-circuit. Narrowed to ErrFailsafeCircuitBreakerOpen
  to avoid double-counting failures that already reach tryForward via the
  normal forwarding path.

* Per-request escape hatch in Network.Forward: when the inner loop
  exhausts upsList with retryable errors and failover.onDefaultsExhausted
  is enabled, replace upsList with the fallback-group upstreams that are
  bootstrapped + CB-closed + method-allowed (cordon ignored - the escape
  is precisely what bypasses it), reset attempted/ConsumedUpstreams/
  ErrorsByUpstream for them, and re-enter the inner loop with the
  selectionPolicy permit check bypassed. First fallback whose gate passes
  serves the request. Client never sees a failure.

Steady-state cost is zero - the escape hatch only triggers on
exhaustion, so healthy primaries return on the first iteration and
fallbacks are never touched. Outage cost is one fallback request per
client request for ~1m until the policy promotes fallbacks into upsList
directly, after which the escape hatch is no longer needed per-request.

Observability:
* erpc_network_fallback_escape_total{project,network,category}
  Counts escape-hatch firings. Sustained non-zero in steady state is a
  config signal (alert candidate); sustained non-zero during an outage
  is expected during the policy-reaction window.

Tests:
* erpc/networks_failover_escape_test.go - 4-upstream test layout (2
  primaries + 2 fallbacks) with a production-style selectionPolicy that
  cordons fallbacks while primaries are healthy. Regression test
  TestFailover_GateSkipsAccumulateErrorRate verifies the recording
  fixes. TestFailover_EscapeHatch sub-tests:
    - EscapesToFallbackOnFirstFailingRequest
    - NoEscapeWhenPrimariesHealthy (verifies zero steady-state cost)
    - NoEscapeWhenFailoverDisabled (operator opt-out)
    - OnlyEscalatesOncePerRequest (no escalation loop)

* upstream/registry_fallback_escape_test.go - unit tests for
  GetFallbackEscapeUpstreams: filter-by-group, ignores-cordon (central
  property), filter-by-IgnoreMethods, no-fallbacks-returns-empty.

End-to-end binary test against four mock RPC servers reproduced the
outage scenario locally: 10/10 healthy requests served by primaries
with 0 escape firings; 20/20 outage requests served by fallback with
20 escape firings and 0 client-facing failures. The pre-existing
TestNetwork_SelectionScenarios/StatePollerContributesToErrorRateWhenNotResamplingExcludedUpstreams
caught an initial Fix 3 over-recording bug (resolved by the CB-open-only
narrowing).

Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
eth_getBalance, eth_getTransactionCount, eth_call, and eth_estimateGas
now default to TranslateLatestTag=false and TranslateFinalizedTag=false
— same pattern eth_getBlockByNumber already uses.

Pre-fix behaviour on instant-finality chains (Pharos, evm:1672):
NormalizeHttpJsonRpc translated "latest" into a concrete hex block N
sourced from the state poller, which lags chain head by one polling
interval. Network.GetFinality then classified N ≤ lowestFinalized as
DataFinalityStateFinalized, and the production wildcard cache policy
("network:*, method:*, finality:finalized, ttl:720h") served the same
historical answer at block N for the whole polling window. Chainlink
saw a stale eth_getTransactionCount("latest"), kept rebroadcasting the
same nonce, the upstream returned TX_REPLAY_ATTACK on every retry, and
clients hit ErrUpstreamsExhausted instead of receiving the current
chain state. Confirmed in prod on 2026-05-28 with ~19 nonces of stale
return; eth_getBalance was similarly affected.

With the tag preserved, ExtractBlockReferenceFromRequest returns
blockRef="latest", GetFinality's non-numeric-tag fast path returns
DataFinalityStateRealtime, the upstream evaluates at execution time,
and ResolveCacheBlockRef on the cache layer keys per the network's
current tip (each tip advance is a fresh cache key, concurrent
"latest" requests within the same tip coalesce on one entry).

eth_getCode is left translated — bytecode rarely changes and a
translated-block cache hit is acceptable.

Cross-upstream divergence at block N+1 vs N is now possible for these
methods (already true for eth_getBlockByNumber since ba3e912) and is
absorbed by the 5s realtime cache TTL.

Tests:
- networks_interpolation_test.go: existing tests that asserted "latest"
  was rewritten to hex on eth_getBalance/eth_call/eth_estimateGas now
  assert the literal tag is preserved; AllMethodsCoverage filters
  flipped per-method.
- networks_finality_interpolation_test.go: rewritten — the test name
  was "PreservesRealtimeFinality across translation", which was only
  meaningful because Finality is cached on first call before
  normalization. With the tag preserved end-to-end, finality is
  Realtime by the natural code path (non-numeric tag) and the body of
  the test reflects that.
- LatestTag_ToHex / FinalizedTag_ToHex / UpstreamSkipping_OnInterpolatedLatest
  switched to eth_getStorageAt / eth_getCode (methods that still
  translate by default).
- New tests: StateMethods_PreserveLatestByDefault verifies the new
  defaults for all four methods, EthGetCode_StillInterpolatesLatest
  locks in the deliberate skip.

Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
…tent success

Some chains reject re-submission of an already-accepted transaction with a
"replay attack" error instead of a standard "already known" message. When the
identical signed tx is broadcast more than once (e.g. parallel/hedged sends),
the first submission lands and the duplicates surface this error. Without
mapping it to the idempotency path, those duplicates are treated as upstream
errors, trip circuit breakers, and bubble up as ErrUpstreamsExhausted even
though the transaction confirmed on-chain.

Map the replay-attack rejection to NonceExceptionReasonAlreadyKnown alongside
the existing duplicate-transaction phrases.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
After rebasing onto upstream/main, several tests referenced APIs the
upstream refactors removed or reshaped. Adapt them to the new model:

- Group field → `tier:fallback` tag (common.TagTierFallback) in upstream
  config literals.
- NewUpstreamsRegistry: drop removed scoreRefreshInterval / ScoringConfig
  args; NewNetwork: pass the new *policy.Engine arg.
- Selection policy: EvalFunctionSource/EvalFunction/selectionPolicyEvaluator
  are gone (replaced by policy.Engine); drive/observe the engine instead.
- Tracker recording calls take a DataFinalityState argument.

Failover-escape expectation updates (upstream's default selection policy
now natively tiers to `tier:fallback` via preferTag, independent of our
per-request escape hatch):

- NoEscapeWhenFailoverDisabled: assert the escape *counter* stays flat
  (our opt-out) rather than expecting ErrUpstreamsExhausted — the request
  may still be served by the fallback tier via native policy routing.
- OnlyEscalatesOncePerRequest: assert the escape fires exactly once
  (no re-escalation loop) rather than expecting exhaustion.

- Remove TestNetworkConfig_SetDefaults_FailoverSkipsAutoSelectionPolicy:
  it verified the `!Failover.Enabled()` auto-policy guard, which is moot
  now that NewDefaultNetworkConfig returns an empty config (upstream's
  default already keeps fallback upstreams eligible).

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
…head-liveness

Root-caused from the 2026-06-12 zkSync (chain 324) incident: deleting an
upstream fullnode pod leaves a half-open TCP connection through the
gateway; eRPC's upstream WS client then wedged permanently — believing
itself connected, delivering zero newHeads, while still handing out
client subscription IDs — until the pod was manually restarted.

Four distinct gaps, each fixed:

1. Dead-peer detection (clients/ws_json_rpc_client.go): pings were
   written but pongs never checked — no read deadline, no pong handler —
   so ReadMessage blocked forever on a black-holed connection (ping
   writes 'succeed' into the proxy/kernel buffer, so the write path
   never errored either). Now: liveness deadline (wsPongWait) armed at
   connect, extended on every pong and data frame; on expiry the
   connection is torn down and re-dialed with the existing backoff.
   Failed ping writes also force-close the connection instead of being
   debug-logged. Callbacks now fire synchronously so adapters observe
   disconnect strictly before reconnect.

2. Resubscribe retry (indexer/adapters/wsupstream): subscribe RPCs ride
   upstream.Forward, i.e. the failsafe circuit breaker, which is
   typically still open at the instant the WS layer reconnects after an
   outage. The old single-shot resubscribe failed once with a warning
   and never retried — no heads until the next disconnect. Now a
   per-connection-epoch retry loop (1s..30s backoff) runs until newHeads
   plus all filters are re-established; disconnect cancels the epoch and
   clears the stale sub ID.

3. Circuit breaker defaults (common/defaults.go): SetDefaults assigned
   SuccessThresholdCapacity twice; the first (200) shadowed the intended
   10, so default-config breakers needed an absurd half-open sample.
   (Production configs setting it explicitly were unaffected.)

4. Loud failure instead of silent sub IDs: newHeads subscribe now
   refuses (HTTP 503, ErrNoLiveSubscriptionSource) when no ingress for
   the network has a live connection + active upstream subscription,
   after a short bootstrap grace wait — so clients fail over instead of
   holding a dead subscription. New observability:
   - erpc_upstream_websocket_connected{project,vendor,network,upstream}
   - erpc_network_subscription_last_head_timestamp_seconds{network}
   - healthcheck networks[].subscriptions {live/total ingresses,
     lastHeadAt, lastHeadAgeSeconds}

Regression tests cover: silent (no close handshake) peer death and
reconnect at the client layer; reconnect + breaker-open retry + heads
resuming at the adapter layer; refusal/grace-wait at the subscription
manager; breaker default values.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
Full-stack reproduction of the 2026-06-12 zkSync incident: real eRPC HTTP
server + WS upstream whose TCP connection is killed with no close
handshake. Asserts eRPC re-dials, re-subscribes with a fresh upstream sub
ID, resumes delivering heads to the already-connected client on its
original subscription ID, and accepts new client subscriptions — all with
no process restart.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
Per review feedback (health reporter excessive vs ping/pong), drop the
observability layer and keep only the fixes required to prevent the
incident:

- revert HealthReporter/IngressHealth/LastHead (indexer), healthcheck
  subscriptions block, newHeads refusal gate + ErrNoLiveSubscriptionSource,
  and the last-head-timestamp metric; delete their tests. Kept: client
  liveness, resubscribe retry loop, breaker default fix, and the
  erpc_upstream_websocket_connected gauge.

Race fixes from review:

- snapshot wsPingInterval/wsPongWait (and resub backoff bounds) into
  per-client/per-adapter fields at construction; goroutines no longer read
  package vars, fixing the -race failure where test cleanup restored them
  while a prior client's pingLoop was still running
- pingLoop: write the ping to an explicit conn and compare-and-close via
  teardownConn, so a ping failure can no longer close a fresh connection
  that readLoop already re-dialed
- subscribeNewHeads/subscribeFilter: check epoch ctx under subsMu before
  committing a sub ID, so a cancelled epoch's in-flight subscribe can't
  unregister the live handler installed by the newer epoch
- Stop marks the adapter stopped under resubMu so a reconnect callback
  racing Stop can't start a new resubscribe epoch

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
…edge)

Root cause of the internal-eRPC "probe/selection wedge" (incident 2026-06-24).

Breaker.Record returned early for OutcomeIgnore BEFORE the HalfOpen branch that
releases the trial permit (halfOpenInflight--). A HalfOpen trial reserves a permit
in TryAcquirePermit; if that trial resolves as ignorable (timeout, cancellation,
soft error — precisely what a transient redis/upstream blip produces), the permit
was never released. After enough such trials halfOpenInflight saturates the trial
capacity, every subsequent TryAcquirePermit in HalfOpen is denied, and the breaker
wedges open indefinitely — failing real traffic AND the selection-recovery probes
(which are breaker-eligible). The upstream can never re-admit, so eRPC serves no
healthy upstream for the chain until the pods are rollout-restarted (which resets
in-memory breaker state).

Evidence: during a wedge the selection-probe error RATIO sits at ~1.0 (every probe
denied) sustained for tens of minutes across multiple chains at once, recovering
within minutes of a restart; onset correlates with redis-haproxy churn.

Fix: on OutcomeIgnore, still release a reserved HalfOpen trial permit (without
counting it as success/failure — the trial was inconclusive). Minimal, behaviour-
preserving for Closed/Open.

Test: breaker_test.go reproduces the leak — without the fix the breaker "wedges
after 0 ignored trials" (halfOpenInflight leaks to 1 and TryAcquirePermit denies);
with the fix, repeated ignored trials never wedge and a later success still closes.

Follow-ups (separate): mark selection-recovery probes breaker-ineligible so a
genuinely-open breaker can't blind its own recovery probe; bounded HalfOpen dwell /
cordon TTL as defence-in-depth.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
fix(failsafe): release half-open permit on ignored outcome (breaker wedge / probe-selection wedge root cause)
… is fatal (backport erpc#973)

Backports upstream erpc commit a04655e (PR erpc#973) verbatim, adapted only
for this fork's attemptRemainingTasks signature.

The auto-retry loop exited permanently when State() returned Fatal, which
happens when ANY single task is fatal (e.g. one misconfigured upstream
with a chainId mismatch). Every transiently-failed sibling task was then
abandoned until process restart. In production this meant an upstream
that blipped during a daemonset rollout (robinhood-mainnet) stayed
unregistered (network n/a) for two days while all its traffic escaped to
a paid 3P fallback.

The loop now stops only when every task is terminal (succeeded or fatal);
a fatal task is terminal for itself only.
…sk defence)

Upstream (erpc#973) still waits unbounded on WaitForTasks each retry
round: one task hung inside its Fn (e.g. a client dial that ignores ctx
and never returns) stays Running forever and blocks the retry loop for
every other task. Bound each round's wait by TaskTimeout so a hung task
only delays a round, never stops retries. Candidate for upstreaming.
…lock-then-wait deadlock)

Stop() held tasksMu while waiting for the auto-retry goroutine, which
itself acquires tasksMu inside attemptRemainingTasks. If the loop was
blocked acquiring the mutex when the cancel landed, Stop waited forever.
Cancel-and-wait now happens before taking the mutex. Present upstream as
well; candidate for upstreaming.
Prevent HTTP "latest"/eth_blockNumber from regressing below a head already
delivered on WS, which trips Chainlink MultiNode FinalizedBlockOutOfSync.

Co-authored-by: Cursor <cursoragent@cursor.com>
Co-authored-by: Cursor <cursoragent@cursor.com>
- Reap Running tasks past TaskTimeout to TimedOut so ignore-ctx Fns cannot
  wedge hasPendingWork forever; completion uses attempt-ID + CAS so a late
  return cannot clobber a newer retry.
- State(): one fatal sibling no longer maps the whole Initializer to Fatal
  while others succeeded/recover — mix is Partial; all-fatal stays Fatal.
- waitForTasks waits in parallel and distinguishes wait-context abort from
  a task that already finished as TimedOut; Wait surfaces TimedOut/Fatal
  errors via lastErr.

Co-authored-by: Cursor <cursoragent@cursor.com>
fix(initializer): keep auto-retry alive for transient tasks despite fatal or hung siblings
…test

fix(ws): advance network latest tip before newHeads fan-out
TipHW alone still allowed eth_getBlockByNumber("latest") to fail-open to a
stale upstream block when re-fetch of the WS tip missed. Cache the newHeads
header before fan-out and prefer it in EnforceHighestBlock for header-only
requests so HTTP cannot regress below a tip already delivered on the pod.

Co-authored-by: Cursor <cursoragent@cursor.com>
Record the tip-source upstream id with the cached newHeads head, partition
eth_getBlockByNumber("latest") using TipHW so WS upstreams are tried first,
and pin EnforceHighestBlock re-fetch to that tip source instead of only
excluding the stale HTTP responder.

Co-authored-by: Cursor <cursoragent@cursor.com>
Tip ownership already lives on per-upstream SuggestLatestBlock/LatestBlock
and partitionUpstreamsByLatestBlock. Keep TipHW partitioning for "latest"
plus the header cache floor; remove the redundant upstreamId pin registry.

Co-authored-by: Cursor <cursoragent@cursor.com>
Drop the newHeads header-cache / tip-source registry overbuild. Tip
ownership already lives on per-upstream pollers (SuggestLatestBlock) and
EvmLeaderUpstream; EnforceHighestBlock now UseUpstream-pins the concrete
tip re-fetch to that leader when its LatestBlock covers TipHW.

Co-authored-by: Cursor <cursoragent@cursor.com>
Availability checks and EnforceHighestBlock were calling PollLatestBlockNumber,
which can reuse a debounced tip behind network TipHW (WS/Redis), falsely
rejecting eth_call and fail-opening stale latest. Add PollLatestBlockNumberNow
and use it on those paths; force-poll the leader before pinning the tip re-fetch.

Co-authored-by: Cursor <cursoragent@cursor.com>
EnforceHighestBlock used pickHighestBlock, which fail-opened to a lagging
"latest" when the concrete TipHW fetch returned null/error. That is the
MultiNode FOOS trigger after WS newHeads already delivered the higher head.

Re-fetch tip (leader pin, then unconstrained), accept only responses that
meet the tip floor, and return an error instead of stale. Also fail-open
the per-upstream availability gate when poller lags TipHW, and only skip
enforcement for cached latest that already meets the tip.

Co-authored-by: Cursor <cursoragent@cursor.com>
Cross-pod TipHW Redis push was async, so sibling pods could still serve
a lower HTTP latest after WS newHeads advanced highestUserObservations,
silently demoting MultiNode via FOOS. Publish TipHW before fan-out and
refresh from Redis when local TipHW would skip EnforceHighestBlock.

Co-authored-by: Cursor <cursoragent@cursor.com>
If TipHW advanced from a WS newHeads observation on this pod, HTTP
eth_getBlockByNumber("latest", false) must return that header instead of
hard-failing when concrete tip re-fetch cannot reach TipHW yet.

Co-authored-by: Cursor <cursoragent@cursor.com>
jleeh and others added 11 commits July 23, 2026 14:43
Fallback newHeads must not advance TipHW while primaries are up (same
rule as poller aggregation). Also let the per-request fallback escape
hatch fire when primaries return null/emptyish for a concrete tip block,
so tip re-fetch can use public fallbacks instead of hard-failing.

Co-authored-by: Cursor <cursoragent@cursor.com>
That path masked TipHW inflation / missing failover. Keep TipHW floor
from primary WS only, refuse-stale when tip re-fetch misses, and escape
to fallbacks on emptyish primary misses.

Co-authored-by: Cursor <cursoragent@cursor.com>
Skipping TipHW for fallback WS while Ingest still fans those heads out
caused MultiNode FOOS on matic. TipHW must cover every delivered head;
tip re-fetch of a fallback-advanced TipHW relies on the emptyish escape
hatch to reach public upstreams when primaries miss.

Co-authored-by: Cursor <cursoragent@cursor.com>
…hed-head

fix(ws): refuse stale latest below TipHW (FOOS)
Direct tip/tip+1 getBlock reads set UseUpstream to the primary leader
(typically the WS ingress advanced by SuggestLatestBlock) on first
forward — same idea as PR11 tip re-fetch pin, so the lagging sibling is
not tried first.

Co-authored-by: Cursor <cursoragent@cursor.com>
Drop helper + dedicated HTTP test; keep only the small UseUpstream pin
next to partitionUpstreamsByLatestBlock.

Co-authored-by: Cursor <cursoragent@cursor.com>
Flatten nested Forward pin into early-return helper; assert lagging
sibling gets zero hits on concrete tip getBlock.

Co-authored-by: Cursor <cursoragent@cursor.com>
…vm-leader

fix(ws): pin near-tip getBlock to EvmLeaderUpstream
Manual publish of experiment/prod tags from an arbitrary git ref
(default feat/websocket-support). Requires DOCKERHUB_USERNAME and
DOCKERHUB_TOKEN repo secrets with write to the linkpool Hub org.

Co-authored-by: Cursor <cursoragent@cursor.com>
Co-authored-by: Cursor <cursoragent@cursor.com>
@github-actions

github-actions Bot commented Sep 15, 2026 •

Copy link
Copy Markdown
File Lines Key changes Risk
🟠 networks.go +110/-0 maxLatestStateLagBlocks, isLatestStateReadMethod, isMovingLatestOrPendingTag, ... ⚠ ErrUpstreamBlockUnavailable
🔵 tracker.go +19/-0 loadOrStoreUpsMetrics
3 test files +381

xray — see through AI slop with deterministic architecture PR diff reviews

@1marcghannam
1marcghannam changed the base branch from main to fix/pin-near-tip-getblock-to-evm-leader September 15, 2026 09:40
@1marcghannam
1marcghannam force-pushed the fix/eth-call-latest-skip-lagging-upstreams branch from d9a3a20 to ac2e98e Compare September 15, 2026 09:41
… lag to */All

Stale eth_call("latest") could be served by a WS peer thousands of blocks behind TipHW because tip integrity only covered getBlock and bn=0 failed open. Skip peers lagging >16 for moving-tag state reads, prefer near-tip order, and mirror BlockHeadLag into {*,All} so network-scope lag predicates see quiet/stalled peers.
…skip

Use the stalled-vs-tip totalQueued values (16355 vs 692), ~8k block lag,
and the real selector so the regression mirrors the production failure mode.
@1marcghannam
1marcghannam force-pushed the fix/eth-call-latest-skip-lagging-upstreams branch from c82e538 to 0c96ef1 Compare September 15, 2026 09:50
@1marcghannam
1marcghannam changed the base branch from fix/pin-near-tip-getblock-to-evm-leader to feat/websocket-support September 15, 2026 09:51
1marcghannam and others added 2 commits September 15, 2026 13:26
…tateLagBlocks

The 16-block hard gate mirrored the default selection policy's
blockNumberLagAbove(16) but was not operator-tunable, unlike the policy
threshold. Add evm.maxLatestStateLagBlocks (nil -> 16, <=0 disables),
resolved at the read site like MaxRetryableBlockDistance. Also correct
the lag-test comments: lag comes from the standard poller mocks
(0x11118888 vs 0x22228888), not the no-op SuggestLatestBlock calls.
@pnpm/exe@10.28.2 (from packageManager) has no linux-x64-musl binary,
which breaks Publish to GHCR on node:alpine. Pin npm-global pnpm and
disable package-manager version switching for install steps.

Co-authored-by: Cursor <cursoragent@cursor.com>
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

4 participants