Skip to content

Use selectors for subscriptions - #401

Merged
KirillPamPam merged 2 commits into
mainfrom
sub_selectors
Oct 6, 2026
Merged

KirillPamPam merged 2 commits into
mainfrom
sub_selectors

Conversation

@KirillPamPam

Copy link
Copy Markdown
Collaborator

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 without eth_getLogs
for logs (no eligible upstream has the block yet, so its logs are skipped). PR #398 patched the
logs 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 upstreams
that pass the filter, with membership re-evaluated on every upstream event. The newHeads and
logs sources 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 a HeightForkChoice and a member set
and 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 Available plus 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 are
    tracked 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(), which
starts 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). WsCapMatcher is generalised to CapMatcher(cap, method), and matcherFilter adapts a Matcher to the supervisor's FilterUpstream.

source feed filter aggregation key
newHeads Available + NewHeadsCap + client selectors local|newHeads|<selectors>
logs Available + LogsCap + client selectors local|logs|<selectors>

LogsCap already implies a subscription-driven head plus eth_getLogs. The logs fetch side
is unchanged: per block, eth_getLogs {blockHash} goes to the best-rated upstream at or above the
height, with no client selectors, since logs fetched by hash are identical on every node (dshackle's
ProduceLogs ignores the selector there too). The #398 failure cannot recur because the head's own
upstream passed the feed filter, so it has the method and the block.

StreamBlockUpdates is event-driven again and takes the feed; parent-hash backfill of skipped
heights and the not-ready retry path from #398 stay.

Resolution (resolveSource). Returns (resolvedSource{key, builder, filter}, error). With a
topic's local-subscriptions flag on, the subscription is served locally or fails with
no available upstreams; there is no silent fallback onto a node the client did not select or whose
head is polled. The flag is the operator's switch for chains that cannot serve the topic locally.

Behaviour changes

  • Local newHeads/logs are strict. With the per-chain flag on (the default), no matching
    upstream 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: false
    and enable-logs: false in chain-defaults, or clients get the error. newPendingTransactions
    keeps its capability-based fallback.
  • Client selectors now apply to local newHeads and logs. They narrow the upstreams whose
    heads the source follows; clients with different selectors get different sources. Previously
    newHeads ignored selectors and logs with selectors always went to a node.
  • A source's first client receives the current head (and the current block's logs), as in
    dshackle; before, the source waited for the next block.
  • RequestAnySelector no longer splits aggregation keys (it has no routing effect), for local and
    generic subscriptions alike.
  • The gauge nodecore_logs_source_head_lag_blocks is removed; it measured the fix(logs-source): follow the heads of upstreams that can serve logs #398 workaround.
    nodecore_logs_source_backfill_failed_total and reorg_clamped_total stay.

Unchanged: the global merged head, SubscribeState, emerald SubscribeChainStatus, head-lag
tracking, ForkChoice.Choose's signature. HeadUpdated.UpstreamId is the upstream whose event
produced the head, the same convention as ChainHeadData.UpstreamId; when a feed's head moves down
because its owner left, the lower head is forwarded like any other, as the global-head source did.

Removed from #398

bestHead, the 50 ms recheck ticker, the per-event full state scan, the one-second lostCheck
ticker, logsMatcher/canServeLogs, the head-lag gauge. Kept: advanceWithAncestors,
BlockResolver, maxBackfillBlocks, the not-ready retry.

Other changes

  • fork_choice.NewHeightForkChoice returns the ForkChoice interface and
    NewGenericChainSupervisor takes a factory, so every feed gets its own instance.
  • NewHeadFeedSubscription exists so ChainSupervisor fakes 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/publish through a
    scripted fork choice.
  • subengine: newHeads source on a fake feed (forwarding, empty at start and mid-stream,
    unsubscribe on cancel, filter pass-through); StreamBlockUpdates harness (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 and
    logs, selectorKey ignoring RequestAnySelector, processor-level terminal frame when nothing can
    serve local newHeads; all existing TestLogsSource* pass unchanged in behaviour.
  • make lint: 0 issues. make test (go test -race ./...): all packages pass.

@KirillPamPam
KirillPamPam merged commit 1bdb753 into main Oct 6, 2026
5 checks passed
@KirillPamPam
KirillPamPam deleted the sub_selectors branch October 6, 2026 09:26
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.

2 participants