diff --git a/docs/pages/operation/websocket.mdx b/docs/pages/operation/websocket.mdx index 3408a61b4..6fd0aa9a4 100644 --- a/docs/pages/operation/websocket.mdx +++ b/docs/pages/operation/websocket.mdx @@ -137,7 +137,7 @@ ws.send(JSON.stringify({ ### Subscription features -- **Resilient delivery** — `newHeads` is subscribed on every connected WebSocket upstream of the network, and each block is delivered from whichever upstream announces it first. `logs` and `newPendingTransactions` are subscribed on every WebSocket upstream whose circuit breaker is closed; with [`failover.onDefaultsExhausted`](/config/projects/networks#failover) enabled, upstreams tagged `tier:fallback` are used only when every other upstream rejects the subscribe. +- **Resilient delivery** — `newHeads` is subscribed on every connected WebSocket upstream of the network, and each block is delivered from whichever upstream announces it first. With [`failover.onDefaultsExhausted`](/config/projects/networks#failover) enabled, heads from upstreams tagged `tier:fallback` are delivered only while no other upstream that the selection policy routes to has a live `newHeads` subscription, so clients are not told of a block only a fallback has while the primaries are up. A held-back head still updates that fallback's tip tracker. How quickly fallback heads take over when the primaries fail follows the selection policy's `evalInterval`. `logs` and `newPendingTransactions` are subscribed on every WebSocket upstream whose circuit breaker is closed; with [`failover.onDefaultsExhausted`](/config/projects/networks#failover) enabled, upstreams tagged `tier:fallback` are used only when every other upstream rejects the subscribe. - **No surprise disconnects** — upstream reconnects, scoring changes and brief outages never close your connection; eRPC re-subscribes upstream behind the scenes. eRPC closes a connection only on shutdown (`1001`), when the client stops answering pings, or when a `logs` subscription overflows its buffer (see below). - **Deduplication** — even though multiple upstreams are subscribed in parallel, you only receive each block, log or pending tx once. Logs are identified by `blockHash`, `transactionHash` and `logIndex`; a log missing any of them is passed through without deduplication. - **Shared subscriptions** — clients subscribing to the same event on the same network share one upstream subscription per upstream. Each client gets its own subscription ID and receives every notification independently. @@ -329,7 +329,7 @@ eRPC minimises the gap in three ways: - **Per-block cache keying for moving tags.** The realtime cache entry for a `latest` / `finalized` request is keyed by the block number the tag currently resolves to, so each tip advance is a distinct key rather than pinning every call within the TTL to the same cached payload. - **`enforceHighestBlock` response rewrite.** When the upstream returns a block number below the network's tracked tip, eRPC's [integrity check](/config/failsafe/integrity) fetches the specific higher block from another upstream and returns that instead (falling back to the highest valid block it has if none can serve it). -- **Delivered-head floor.** When an upstream WebSocket delivers a `newHeads` notification, eRPC updates that upstream's tip tracker before fanning the head out to clients, and records it as a floor for the network's `latest`. Every head a subscriber can receive feeds the floor once the delivering upstream's tip tracker has accepted it, and the floor is bounded: a value no live upstream comes close to (a stale or poisoned entry) is ignored rather than trusted, so it cannot pin the served tip. It never lifts `latest` past [`servedTip.guaranteedFor`](/config/projects/networks) / `guaranteedMethods`, and it does not apply to requests scoped with `use-upstream`. With shared state configured the floor is also propagated to other instances in the background, so they converge within the propagation delay rather than instantly. +- **Delivered-head floor.** When an upstream WebSocket delivers a `newHeads` notification, eRPC updates that upstream's tip tracker before fanning the head out to clients, and records it as a floor for the network's `latest`. Every head a subscriber can receive (so not a held-back fallback head) feeds the floor once the delivering upstream's tip tracker has accepted it, and the floor is bounded: a value no live upstream comes close to (a stale or poisoned entry) is ignored rather than trusted, so it cannot pin the served tip. It never lifts `latest` past [`servedTip.guaranteedFor`](/config/projects/networks) / `guaranteedMethods`, and it does not apply to requests scoped with `use-upstream`. With shared state configured the floor is also propagated to other instances in the background, so they converge within the propagation delay rather than instantly. Even so, a transient single-block lag can reach a client when the chain produces a block faster than an in-flight HTTP round-trip completes. If your consumer treats single-block HTTP-vs-WS lag as a hard failure, widen its tolerance to match the chain's block time rather than expect eRPC (or any proxy) to eliminate it. diff --git a/erpc/network_executor.go b/erpc/network_executor.go index cab3dd46d..8be2b55b2 100644 --- a/erpc/network_executor.go +++ b/erpc/network_executor.go @@ -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. diff --git a/erpc/networks.go b/erpc/networks.go index be977e951..d1ad6ebd6 100644 --- a/erpc/networks.go +++ b/erpc/networks.go @@ -3627,9 +3627,12 @@ func (n *Network) acquireRateLimitPermit(ctx context.Context, req *common.Normal // tierUpstreamsByGroup moves fallback-tier upstreams behind the rest, // preserving order within each tier. func tierUpstreamsByGroup(ups []common.Upstream) []common.Upstream { - return stablePartition(ups, func(u common.Upstream) bool { - return u.Config() != nil && u.Config().HasTag(common.TagTierFallback) - }) + return stablePartition(ups, isFallbackTier) +} + +// isFallbackTier reports whether u is tagged tier:fallback. +func isFallbackTier(u common.Upstream) bool { + return u.Config() != nil && u.Config().HasTag(common.TagTierFallback) } // partitionUpstreamsByLatestBlock moves upstreams whose polled head is known diff --git a/erpc/networks_failover_escape_test.go b/erpc/networks_failover_escape_test.go index 1f9f60cac..97434ac63 100644 --- a/erpc/networks_failover_escape_test.go +++ b/erpc/networks_failover_escape_test.go @@ -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 { @@ -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( @@ -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) diff --git a/erpc/networks_failover_hedge_test.go b/erpc/networks_failover_hedge_test.go new file mode 100644 index 000000000..cf6d1765f --- /dev/null +++ b/erpc/networks_failover_hedge_test.go @@ -0,0 +1,167 @@ +package erpc + +import ( + "context" + "net/http" + "strings" + "sync/atomic" + "testing" + "time" + + "github.com/erpc/erpc/common" + "github.com/erpc/erpc/telemetry" + "github.com/erpc/erpc/util" + "github.com/h2non/gock" + promUtil "github.com/prometheus/client_golang/prometheus/testutil" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +// unboundedAll drops the availability bound from every upstream, so each is +// called (and answers) for blocks above its polled head. +func unboundedAll(cfgs []*common.UpstreamConfig) { + for _, cfg := range cfgs { + cfg.Evm.BlockAvailability = nil + } +} + +// ethCallMock answers eth_call on host with body after delay. +func ethCallMock(host string, delay time.Duration, body string) { + 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). + Delay(delay). + JSON([]byte(body)) +} + +// countEthCalls counts eth_call requests reaching host, ahead of its standard +// mock. +func countEthCalls(host string, hits *atomic.Int64) { + gock.New("http://" + host). + Post(""). + Persist(). + Filter(func(r *http.Request) bool { + // Filters run before host matching. + if r.URL.Host == host && strings.Contains(util.SafeReadBody(r), "eth_call") { + hits.Add(1) + } + return false + }). + Reply(200) +} + +const missingDataBody = `{"jsonrpc":"2.0","id":1,"error":{"code":-32000,"message":"header not found"}}` + +// hedgeFaster hedges well before a slow upstream answers. +func hedgeFaster() []*common.FailsafeConfig { + return []*common.FailsafeConfig{{ + MatchMethod: "*", + Hedge: &common.HedgePolicyConfig{Delay: common.NewStaticDuration(20 * time.Millisecond), MaxCount: 1}, + }} +} + +func forwardEthCall(t *testing.T, ctx context.Context, network *Network, id int, blockHex string) (string, error) { + t.Helper() + req := ethCallRequest(id, blockHex) + req.SetNetwork(network) + resp, err := network.Forward(ctx, req) + if err != nil { + return "", err + } + require.NotNil(t, resp) + defer resp.Release() + jrr, err := resp.JsonRpcResponse() + require.NoError(t, err) + return strings.Trim(jrr.GetResultString(), `"`), nil +} + +// A hedge leg that finds every routed upstream missing the data must not +// cancel a sibling leg that escalated to the fallbacks and is still waiting. +func TestFailover_HedgeKeepsEscalatedSibling(t *testing.T) { + defer util.ResetGock() + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + + // The primaries miss, the escape sends one leg to the (slow) fallbacks, + // and the hedge leg misses on the primaries again. + network, _, _ := setupFailoverFixture(t, ctx, failoverFixtureOpts{ + primaryLatest: "0x3e8", // 1000 + fallbackLatest: "0x3e8", // 1000 + enableFailover: true, + configure: unboundedAll, + failsafe: hedgeFaster(), + mocks: func() { + ethCallMock("rpc1.localhost", 0, missingDataBody) + ethCallMock("rpc2.localhost", 0, missingDataBody) + ethCallMock("rpc3.localhost", 80*time.Millisecond, `{"jsonrpc":"2.0","id":1,"result":"0x3333"}`) + ethCallMock("rpc4.localhost", 80*time.Millisecond, `{"jsonrpc":"2.0","id":1,"result":"0x3333"}`) + }, + }) + + escape := telemetry.MetricNetworkFallbackEscapeTotal.WithLabelValues("main", "evm:999", "eth_call") + escapeBefore := promUtil.ToFloat64(escape) + for i := 0; i < 5; i++ { + result, err := forwardEthCall(t, ctx, network, i, "0x3ea") + require.NoError(t, err, "iter %d", i) + assert.Equal(t, "0x3333", result, "iter %d", i) + } + assert.Equal(t, escapeBefore+5, promUtil.ToFloat64(escape)) +} + +// A fallback that already answered missing data must not let a hedge leg +// cancel another fallback that is still working on the request. +func TestFailover_HedgeKeepsSlowFallbackAfterFastFallbackMiss(t *testing.T) { + defer util.ResetGock() + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + + network, _, _ := setupFailoverFixture(t, ctx, failoverFixtureOpts{ + primaryLatest: "0x3e8", // 1000 + fallbackLatest: "0x3e8", // 1000 + enableFailover: true, + configure: unboundedAll, + failsafe: hedgeFaster(), + mocks: func() { + ethCallMock("rpc1.localhost", 0, missingDataBody) + ethCallMock("rpc2.localhost", 0, missingDataBody) + ethCallMock("rpc3.localhost", 0, missingDataBody) + ethCallMock("rpc4.localhost", 80*time.Millisecond, `{"jsonrpc":"2.0","id":1,"result":"0x4444"}`) + }, + }) + + for i := 0; i < 5; i++ { + result, err := forwardEthCall(t, ctx, network, i, "0x3ea") + require.NoError(t, err, "iter %d", i) + assert.Equal(t, "0x4444", result, "iter %d", i) + } +} + +// A fallback whose polled head is ahead does not take reads the primaries +// serve: fallbacks are reached only through the escape. (The default +// policy's probeExcluded may still mirror the request to them in the +// background, so upstream hits are not counted here.) +func TestFailover_FallbackAheadDoesNotTakeServedReads(t *testing.T) { + defer util.ResetGock() + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + + network, _, _ := setupFailoverFixture(t, ctx, failoverFixtureOpts{ + primaryLatest: "0x3e8", // 1000 + fallbackLatest: "0x3ea", // 1002 + enableFailover: true, + configure: unboundedAll, + }) + + escape := telemetry.MetricNetworkFallbackEscapeTotal.WithLabelValues("main", "evm:999", "eth_call") + escapeBefore := promUtil.ToFloat64(escape) + for i := 0; i < 5; i++ { + result, err := forwardEthCall(t, ctx, network, i, "0x3ea") + require.NoError(t, err, "iter %d", i) + assert.Contains(t, []string{"0x1111", "0x2222"}, result, "iter %d", i) + } + assert.Equal(t, escapeBefore, promUtil.ToFloat64(escape)) +} diff --git a/erpc/networks_ws_tip_test.go b/erpc/networks_ws_tip_test.go index d86d83723..bd3c69cf1 100644 --- a/erpc/networks_ws_tip_test.go +++ b/erpc/networks_ws_tip_test.go @@ -10,6 +10,7 @@ import ( "github.com/erpc/erpc/common" "github.com/erpc/erpc/data" "github.com/erpc/erpc/health" + "github.com/erpc/erpc/internal/policy" "github.com/erpc/erpc/thirdparty" "github.com/erpc/erpc/upstream" "github.com/erpc/erpc/util" @@ -154,6 +155,104 @@ func TestNetworkHandle_SuggestLatestBlock_EveryKnownSourceLiftsLatest(t *testing "a head a client can receive from the cordoned fallback floors latest") } +type fakeHeadSource bool + +func (f fakeHeadSource) HeadsLive() bool { return bool(f) } + +// With failover on, a fallback's head reaches clients (and lifts latest) only +// when no primary the policy routes to is streaming heads itself. The +// fallback's own poller sees every head regardless. +func TestNetworkHandle_FallbackHeadsHeldWhilePrimariesStream(t *testing.T) { + setup := func(t *testing.T, ctx context.Context) (*Network, map[string]*upstream.Upstream) { + network, ups, _ := setupFailoverFixture(t, ctx, failoverFixtureOpts{ + primaryLatest: "0x3e8", // 1000 + fallbackLatest: "0x3e8", + enableFailover: true, + }) + require.Equal(t, int64(1000), network.EvmHighestLatestBlockNumber(ctx)) + byId := make(map[string]*upstream.Upstream, len(ups)) + for _, u := range ups { + byId[u.Id()] = u + } + return network, byId + } + + t.Run("HeldWhileAPrimaryStreams", func(t *testing.T) { + defer util.ResetGock() + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + network, ups := setup(t, ctx) + + handle := &networkHandle{nw: network, heads: map[string]headSource{ + "primary-1": fakeHeadSource(false), + "primary-2": fakeHeadSource(true), + "fallback-1": fakeHeadSource(true), + }} + assert.False(t, handle.SuggestLatestBlock("ws:fallback-1", 1010)) + assert.Equal(t, int64(1010), ups["fallback-1"].EvmStatePoller().LatestBlock(), + "the fallback's poller still sees its head") + assert.Equal(t, int64(1000), network.EvmHighestLatestBlockNumber(ctx), + "a held head must not lift latest") + + assert.True(t, handle.SuggestLatestBlock("ws:primary-2", 1001)) + assert.Equal(t, int64(1001), network.EvmHighestLatestBlockNumber(ctx)) + }) + + t.Run("DeliveredWhenNoPrimaryStreams", func(t *testing.T) { + defer util.ResetGock() + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + network, _ := setup(t, ctx) + + // Primaries connected over HTTP only, or with no live subscription. + for _, heads := range []map[string]headSource{ + {"fallback-1": fakeHeadSource(true)}, + {"primary-1": fakeHeadSource(false), "fallback-1": fakeHeadSource(true)}, + } { + handle := &networkHandle{nw: network, heads: heads} + assert.True(t, handle.SuggestLatestBlock("ws:fallback-1", 1010)) + } + assert.Equal(t, int64(1010), network.EvmHighestLatestBlockNumber(ctx)) + }) + + t.Run("DeliveredWhenPolicyExcludesPrimaries", func(t *testing.T) { + defer util.ResetGock() + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + network, ups := setup(t, ctx) + + handle := &networkHandle{nw: network, heads: map[string]headSource{ + "primary-1": fakeHeadSource(true), + "primary-2": fakeHeadSource(true), + "fallback-1": fakeHeadSource(true), + }} + require.False(t, handle.SuggestLatestBlock("ws:fallback-1", 1010)) + + ups["primary-1"].Cordon("*", "test") + ups["primary-2"].Cordon("*", "test") + policy.TickForTest(network.policyEngine, network.networkId, "*") + + assert.True(t, handle.SuggestLatestBlock("ws:fallback-1", 1011)) + assert.Equal(t, int64(1011), network.EvmHighestLatestBlockNumber(ctx)) + }) + + t.Run("DeliveredWithFailoverOff", func(t *testing.T) { + defer util.ResetGock() + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + network, _, _ := setupFailoverFixture(t, ctx, failoverFixtureOpts{ + primaryLatest: "0x3e8", + fallbackLatest: "0x3e8", + }) + + handle := &networkHandle{nw: network, heads: map[string]headSource{ + "primary-1": fakeHeadSource(true), + "fallback-1": fakeHeadSource(true), + }} + assert.True(t, handle.SuggestLatestBlock("ws:fallback-1", 1010)) + }) +} + // The delivered-head floor is network-wide: a use-upstream-scoped request is // answered from its group's own head. func TestDeliveredHeadFloor_SkipsSelectorScopedRequests(t *testing.T) { diff --git a/erpc/subscription_manager.go b/erpc/subscription_manager.go index b4e0ac8d2..251e185b8 100644 --- a/erpc/subscription_manager.go +++ b/erpc/subscription_manager.go @@ -259,14 +259,22 @@ func (sm *SubscriptionManager) bootstrapNetwork(ctx context.Context, nw *Network return common.NewErrNoWsUpstreamAvailable(networkID) } - sm.idx.RegisterNetwork(&networkHandle{nw: nw}) - sm.idx.RegisterNetworkSelector(networkID, &subIngressSelector{nw: nw}) var adapterOpts wsupstream.Options if cfg := nw.cfg; cfg != nil && cfg.Evm != nil && cfg.Evm.StripSubscribeFromBlockZero != nil { adapterOpts.StripSubscribeFromBlockZero = *cfg.Evm.StripSubscribeFromBlockZero } + adapters := make(map[string]*wsupstream.Adapter, len(wsUpstreams)) + heads := make(map[string]headSource, len(wsUpstreams)) for _, up := range wsUpstreams { - adapter := wsupstream.New(up, networkID, sm.logger, adapterOpts) + if adapter := wsupstream.New(up, networkID, sm.logger, adapterOpts); adapter != nil { + adapters[up.Id()] = adapter + heads[up.Id()] = adapter + } + } + sm.idx.RegisterNetwork(&networkHandle{nw: nw, heads: heads}) + sm.idx.RegisterNetworkSelector(networkID, &subIngressSelector{nw: nw}) + for _, up := range wsUpstreams { + adapter := adapters[up.Id()] if adapter == nil { continue } @@ -416,38 +424,65 @@ func (sm *SubscriptionManager) recordFailureMetrics( // networkHandle adapts *Network to indexer.NetworkHandle. type networkHandle struct { nw *Network + // heads are the network's WebSocket ingresses by upstream id; fixed once + // the network is registered. + heads map[string]headSource +} + +// headSource is an ingress that streams newHeads. +type headSource interface { + // HeadsLive reports whether its newHeads subscription is live. + HeadsLive() bool } func (h *networkHandle) Id() string { return h.nw.networkId } // SuggestLatestBlock passes a head from the ingress "ws:" to -// that upstream's state poller. Once the poller has accepted it (a major -// jump is verified asynchronously first), a head from a tip candidate also -// advances the network's delivered-head floor before clients see it; a -// fallback-tier or policy-excluded upstream cannot lift "latest". -func (h *networkHandle) SuggestLatestBlock(sourceId string, blockNumber int64) { +// that upstream's state poller and reports whether clients may receive it +// (see deliversHeadsFrom). Once the poller has accepted a deliverable head (a +// major jump is verified asynchronously first), it also advances the +// network's delivered-head floor before clients see it. +func (h *networkHandle) SuggestLatestBlock(sourceId string, blockNumber int64) bool { upstreamID, ok := strings.CutPrefix(sourceId, "ws:") if !ok { - return + return true } ctx := context.Background() for _, u := range h.nw.upstreamsRegistry.GetNetworkUpstreams(ctx, h.nw.networkId) { if u.Id() != upstreamID { continue } + deliver := h.deliversHeadsFrom(ctx, u) poller := u.EvmStatePoller() if poller == nil || poller.IsObjectNull() { - return + return deliver } poller.SuggestLatestBlock(blockNumber) - // Every source's heads reach clients (whichever delivers first), so - // every source feeds the floor — once its poller accepted the head, - // which excludes jumps still pending chain-id verification. - if poller.LatestBlock() >= blockNumber { + if deliver && poller.LatestBlock() >= blockNumber { h.nw.NoteObservedLatestBlock(h.nw.appCtx, blockNumber) } - return + return deliver + } + return true +} + +// deliversHeadsFrom reports whether clients may receive u's heads. With +// failover on, a fallback-tier upstream's heads are held back while some +// upstream outside that tier, which the selection policy routes to, has a +// live newHeads subscription of its own: clients are not told of a block +// only a fallback has while the primaries are up. Whether they are up is the +// policy's verdict, as for reads. Without such a primary (none eligible, or +// none streaming heads) the fallbacks' heads are delivered. +func (h *networkHandle) deliversHeadsFrom(ctx context.Context, u common.Upstream) bool { + if h.nw.cfg == nil || !h.nw.cfg.Failover.Enabled() || !isFallbackTier(u) { + return true + } + for _, c := range h.nw.tipCandidateUpstreams(ctx, "*") { + if s := h.heads[c.Id()]; s != nil && !isFallbackTier(c) && s.HeadsLive() { + return false + } } + return true } var ( diff --git a/indexer/adapters/wsupstream/adapter.go b/indexer/adapters/wsupstream/adapter.go index d0814ea9b..0f3caa545 100644 --- a/indexer/adapters/wsupstream/adapter.go +++ b/indexer/adapters/wsupstream/adapter.go @@ -116,6 +116,14 @@ func New(up *upstream.Upstream, networkID string, logger *zerolog.Logger, opts O func (a *Adapter) Name() string { return "ws:" + a.upstreamID } +// HeadsLive reports whether the newHeads subscription is live on the +// current connection. +func (a *Adapter) HeadsLive() bool { + a.subsMu.Lock() + defer a.subsMu.Unlock() + return a.heads.id != "" +} + // Start registers the connection hooks and subscribes newHeads in the // background, detached from ctx. func (a *Adapter) Start(_ context.Context, nw indexer.NetworkHandle, sink indexer.Sink) error { diff --git a/indexer/adapters/wsupstream/adapter_reconnect_test.go b/indexer/adapters/wsupstream/adapter_reconnect_test.go index 672e152c0..ec46a84c7 100644 --- a/indexer/adapters/wsupstream/adapter_reconnect_test.go +++ b/indexer/adapters/wsupstream/adapter_reconnect_test.go @@ -26,8 +26,8 @@ import ( type fakeNetworkHandle struct{} -func (fakeNetworkHandle) Id() string { return "evm:123" } -func (fakeNetworkHandle) SuggestLatestBlock(string, int64) {} +func (fakeNetworkHandle) Id() string { return "evm:123" } +func (fakeNetworkHandle) SuggestLatestBlock(string, int64) bool { return true } type fakeSink struct { events chan indexer.StreamEvent @@ -235,10 +235,12 @@ func TestAdapterResubscribesWithRetryAfterReconnect(t *testing.T) { return send(ctx, nq, bypass) } + require.False(t, a.HeadsLive(), "no heads before the subscribe succeeds") require.NoError(t, a.Start(context.Background(), fakeNetworkHandle{}, sink)) require.Eventually(t, subscribedHeads(a), 3*time.Second, 10*time.Millisecond, "adapter never subscribed newHeads despite the failures stopping after %d", failuresPerEpoch) + require.True(t, a.HeadsLive()) require.GreaterOrEqual(t, forwardCalls.Load(), int64(failuresPerEpoch+1), "expected the subscribe to be retried through failures") diff --git a/indexer/indexer.go b/indexer/indexer.go index e382a5d88..11a592f3e 100644 --- a/indexer/indexer.go +++ b/indexer/indexer.go @@ -303,7 +303,9 @@ func (i *Indexer) Ingest(ev StreamEvent) { // Before dedup, so every source's head counts even if another source // already delivered it. if ev.Kind == KindNewHead && !ev.Block.Zero() && ev.SourceId != "" { - ns.handle.SuggestLatestBlock(ev.SourceId, ev.Block.Number) + if !ns.handle.SuggestLatestBlock(ev.SourceId, ev.Block.Number) { + return + } } // The upstream's removed flag is trusted as-is. diff --git a/indexer/indexer_test.go b/indexer/indexer_test.go index fcbe50949..3d0866c88 100644 --- a/indexer/indexer_test.go +++ b/indexer/indexer_test.go @@ -18,7 +18,8 @@ type fakeNetwork struct { id string mu sync.Mutex - suggestedBy map[string][]int64 // sourceId -> block nums seen + suggestedBy map[string][]int64 // sourceId -> block nums seen + held map[string]struct{} // sourceIds whose heads are not delivered } func newFakeNetwork(id string) *fakeNetwork { @@ -29,10 +30,12 @@ func newFakeNetwork(id string) *fakeNetwork { } func (n *fakeNetwork) Id() string { return n.id } -func (n *fakeNetwork) SuggestLatestBlock(sourceId string, block int64) { +func (n *fakeNetwork) SuggestLatestBlock(sourceId string, block int64) bool { n.mu.Lock() + defer n.mu.Unlock() n.suggestedBy[sourceId] = append(n.suggestedBy[sourceId], block) - n.mu.Unlock() + _, held := n.held[sourceId] + return !held } type fakeEgress struct { @@ -271,6 +274,31 @@ func TestIndexer_NewHead_StalerDroppedKeepsStatePollerFed(t *testing.T) { } } +// A head the network does not deliver is still observed, and does not +// advance the dedup marker past the heads the other sources deliver. +func TestIndexer_NewHead_HeldSourceObservedNotDelivered(t *testing.T) { + idx := newIndexer(t) + nw := newFakeNetwork("evm:1") + nw.held = map[string]struct{}{"ws:fb": {}} + idx.RegisterNetwork(nw) + eg := &fakeEgress{name: "eg1", acceptAllHeads: true} + idx.Attach(eg) + + idx.Ingest(StreamEvent{Kind: KindNewHead, NetworkId: "evm:1", SourceId: "ws:fb", Block: BlockRef{Number: 101, Hash: "0xBBB"}}) + if got := eg.count(); got != 0 { + t.Fatalf("held head delivered: got %d", got) + } + idx.Ingest(StreamEvent{Kind: KindNewHead, NetworkId: "evm:1", SourceId: "ws:up1", Block: BlockRef{Number: 101, Hash: "0xBBB"}}) + if got := eg.count(); got != 1 { + t.Fatalf("the delivering source's head must go out: want 1, got %d", got) + } + nw.mu.Lock() + defer nw.mu.Unlock() + if len(nw.suggestedBy["ws:fb"]) != 1 { + t.Fatalf("held source must still be observed, got %v", nw.suggestedBy) + } +} + func TestIndexer_Log_RefcountFanOutAndTeardown(t *testing.T) { idx := newIndexer(t) nw := newFakeNetwork("evm:1") diff --git a/indexer/ingress.go b/indexer/ingress.go index 3df2d4086..26338e672 100644 --- a/indexer/ingress.go +++ b/indexer/ingress.go @@ -11,10 +11,11 @@ type Sink interface { // NetworkHandle is the part of a network an ingress may touch. type NetworkHandle interface { Id() string - // SuggestLatestBlock reports a head observed by source. It is called - // before dedup and fan-out, so every source's observation counts and - // the network tip moves before clients see the head. - SuggestLatestBlock(sourceId string, blockNumber int64) + // SuggestLatestBlock reports a head observed by source and whether + // clients may receive it. It is called before dedup and fan-out, so + // every source's observation counts and the network tip moves before + // clients see the head; a head it does not deliver is dropped. + SuggestLatestBlock(sourceId string, blockNumber int64) (deliver bool) } // EventIngress turns a transport-specific subscription into StreamEvents diff --git a/indexer/integration_test.go b/indexer/integration_test.go index 976144854..ebecd78a8 100644 --- a/indexer/integration_test.go +++ b/indexer/integration_test.go @@ -66,13 +66,14 @@ type stubNetwork struct { } func (s *stubNetwork) Id() string { return s.id } -func (s *stubNetwork) SuggestLatestBlock(sourceID string, block int64) { +func (s *stubNetwork) SuggestLatestBlock(sourceID string, block int64) bool { s.mu.Lock() if s.suggestions == nil { s.suggestions = make(map[string][]int64) } s.suggestions[sourceID] = append(s.suggestions[sourceID], block) s.mu.Unlock() + return true } type recordingEgress struct {