Repository navigation
Use selectors for subscriptions - #401
Merged
Merged
Conversation
KirillPamPam
force-pushed
the
sub_selectors
branch
from
October 5, 2026 16:52
b016f3c to
a13198e
Compare
l0gun0v
approved these changes
Oct 5, 2026
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.
Filtered head feeds: local newHeads/logs follow only the upstreams that can serve them
Summary
The chain supervisor had exactly one merged head per chain, and both locally synthesized EVM
subscriptions tapped it. That head can come from an upstream the source cannot use: a polled-head
upstream for
newHeads(no notification payload to forward), or an upstream withouteth_getLogsfor
logs(no eligible upstream has the block yet, so its logs are skipped). PR #398 patched thelogs case from the consumer side by re-reading every upstream state on each event and every 50 ms.
This PR adds the missing primitive instead, modelled on dshackle's
getHead(matcher):ChainSupervisor.SubscribeHead(name, filter)returns a head feed computed over only the upstreamsthat pass the filter, with membership re-evaluated on every upstream event. The
newHeadsandlogssources are rebuilt on it, keyed per client selector set, and #398's polling is removed.Design
Feed (
internal/upstreams/head_feed.go). Each feed owns aHeightForkChoiceand a member setand is driven from the supervisor's single event loop: after every upstream event the filter is
re-run for that upstream, so method bans, capability loss on a websocket disconnect, status flips,
removal and recovery take effect on the next event, and height-dependent selectors work. Admission
is
Availableplus the filter. The feed emits a sum type:HeadUpdated{Head, UpstreamId}when the fork choice's head changes to a non-empty block;HeadFeedEmpty{}when the member set goes from non-empty to empty. Members and heads aretracked separately so "members exist but none has reported a head yet" is not mistaken for
"nothing passes the filter".
On subscribe the feed seeds its current state, matching dshackle's
AbstractHead.getFlux(), whichstarts with the current head. Unlike dshackle, membership is dynamic rather than a snapshot at
creation time. A subscriber that stops reading (a full 100-event buffer) has its feed closed, the
same way the subscription engine disconnects a too-slow client; nothing is ever dropped.
Consumers (
flow,flow/subengine).WsCapMatcheris generalised toCapMatcher(cap, method), andmatcherFilteradapts aMatcherto the supervisor'sFilterUpstream.newHeadsNewHeadsCap+ client selectorslocal|newHeads|<selectors>logsLogsCap+ client selectorslocal|logs|<selectors>LogsCapalready implies a subscription-driven head pluseth_getLogs. The logs fetch sideis unchanged: per block,
eth_getLogs {blockHash}goes to the best-rated upstream at or above theheight, with no client selectors, since logs fetched by hash are identical on every node (dshackle's
ProduceLogsignores the selector there too). The #398 failure cannot recur because the head's ownupstream passed the feed filter, so it has the method and the block.
StreamBlockUpdatesis event-driven again and takes the feed; parent-hash backfill of skippedheights and the not-ready retry path from #398 stay.
Resolution (
resolveSource). Returns(resolvedSource{key, builder, filter}, error). With atopic's
local-subscriptionsflag on, the subscription is served locally or fails withno available upstreams; there is no silent fallback onto a node the client did not select or whosehead is polled. The flag is the operator's switch for chains that cannot serve the topic locally.
Behaviour changes
newHeads/logsare strict. With the per-chain flag on (the default), no matchingupstream right now is a terminal error for the client instead of a node-backed passthrough. An
EVM chain whose upstreams all have a JSON-RPC head connector must set
enable-new-heads: falseand
enable-logs: falseinchain-defaults, or clients get the error.newPendingTransactionskeeps its capability-based fallback.
newHeadsandlogs. They narrow the upstreams whoseheads the source follows; clients with different selectors get different sources. Previously
newHeadsignored selectors andlogswith selectors always went to a node.dshackle; before, the source waited for the next block.
RequestAnySelectorno longer splits aggregation keys (it has no routing effect), for local andgeneric subscriptions alike.
nodecore_logs_source_head_lag_blocksis removed; it measured the fix(logs-source): follow the heads of upstreams that can serve logs #398 workaround.nodecore_logs_source_backfill_failed_totalandreorg_clamped_totalstay.Unchanged: the global merged head,
SubscribeState, emeraldSubscribeChainStatus, head-lagtracking,
ForkChoice.Choose's signature.HeadUpdated.UpstreamIdis the upstream whose eventproduced the head, the same convention as
ChainHeadData.UpstreamId; when a feed's head moves downbecause its owner left, the lower head is forwarded like any other, as the global-head source did.
Removed from #398
bestHead, the 50 msrecheckticker, the per-event full state scan, the one-secondlostCheckticker,
logsMatcher/canServeLogs, the head-lag gauge. Kept:advanceWithAncestors,BlockResolver,maxBackfillBlocks, the not-ready retry.Other changes
fork_choice.NewHeightForkChoicereturns theForkChoiceinterface andNewGenericChainSupervisortakes a factory, so every feed gets its own instance.NewHeadFeedSubscriptionexists soChainSupervisorfakes in other packages can hand out feeds.Testing
internal/upstreams/chain_supervisor_head_feed_test.go(13 tests against the real supervisor):seeding (current head, empty, and no event while members have no head), admission and eviction on
each event type, height-dependent filter, two feeds over one event stream, unsubscribe
idempotence, shutdown, slow-subscriber close.
internal/upstreams/head_feed_internal_test.go(11 tests):admit/apply/publishthrough ascripted fork choice.
subengine: newHeads source on a fake feed (forwarding, empty at start and mid-stream,unsubscribe on cancel, filter pass-through);
StreamBlockUpdatesharness (follow, close on empty,close on feed close and cancel, backfill via resolver, gap on resolver failure).
flow:CapMatcher,matcherFilter, per-selector keys and strict resolution for newHeads andlogs,
selectorKeyignoringRequestAnySelector, processor-level terminal frame when nothing canserve local newHeads; all existing
TestLogsSource*pass unchanged in behaviour.make lint: 0 issues.make test(go test -race ./...): all packages pass.