fix(evm): skip lagging upstreams for eth_call(latest) and mirror head lag to */All - #25
Open
1marcghannam wants to merge 45 commits into
Open
1marcghannam wants to merge 45 commits into
1marcghannam wants to merge 45 commits into
Conversation
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>
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>
Co-authored-by: Cursor <cursoragent@cursor.com>
xray — see through AI slop with deterministic architecture PR diff reviews |
1marcghannam
changed the base branch from
main
to
fix/pin-near-tip-getblock-to-evm-leader
September 15, 2026 09:40
1marcghannam
force-pushed
the
fix/eth-call-latest-skip-lagging-upstreams
branch
from
September 15, 2026 09:41
d9a3a20 to
ac2e98e
Compare
… 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
force-pushed
the
fix/eth-call-latest-skip-lagging-upstreams
branch
from
September 15, 2026 09:50
c82e538 to
0c96ef1
Compare
1marcghannam
changed the base branch from
fix/pin-near-tip-getblock-to-evm-leader
to
feat/websocket-support
September 15, 2026 09:51
…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>
jleeh
force-pushed
the
feat/websocket-support
branch
from
September 28, 2026 12:31
4c7a584 to
a9e0685
Compare
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
Summary
totalQueuedvia eRPC returned 16,355 instead of ~692 because a stalled WS peer (internal-eth-mainnet-reth-ws-0) still answeredeth_call("latest"). Tip-pin image only covers getBlock.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).latest/pending/empty tag.{ups,"*",All}.Base:
fix/pin-near-tip-getblock-to-evm-leader(LinkPool tip-pin / WS lineage — not upstreammain).Test plan
TestNetwork_EthCallLatest_SkipsLaggingUpstreamTestNetwork_EthCallConcreteBlock_AllowsLaggingArchiveTestWildcardLagMirroredForPeerUpstreamsRollout
internal-erpcimage tag.totalQueuedvia eRPC ≈ 692; no stalled reth-0 while behind TipHW.