diff --git a/benchmarks/aggregator-head-lag.yml b/benchmarks/aggregator-head-lag.yml index 88404b44e..a9ff0d9f6 100644 --- a/benchmarks/aggregator-head-lag.yml +++ b/benchmarks/aggregator-head-lag.yml @@ -70,7 +70,7 @@ methodology: - "Solana (since 2026-09-16): the headline is the lag behind the first feed to report the trade, our own node subscription included in the race. Solana has no on-chain timestamp with sub-second precision, so the timestamps providers send compare conventions: Mobula's `date` (its ingestion time) read as a constant 0.10 s, Serialized's `at` (blockTime, whole seconds) as 0.75 s, while against a common clock the two feeds are 10 to 30 ms apart." - "Solana zero point: no RPC WebSocket we hold, public or keyed (mainnet-beta, Helius, Alchemy), precedes the providers' geyser feeds, so an absolute zero point is not available and the first observation is the ruler. Our node stays in the race: it validates the hash on chain and bounds the field whenever a feed is slower than a plain RPC subscription." - "Race: every feed's arrival for a transaction is compared to the earliest feed's (ties within 5 ms count as first for each feed involved; trades only one feed reported are not scored). Published on every chain as `head_lag_first_share_pct`, the First to report view: the share of the last 24 h of trades a feed reported before every other feed." - - "Race eligibility (since 2026-10-04): only trades our own pool subscription saw are raced. OKX and Birdeye are subscribed per token and see every pool on it, so before this the book was dominated by trades from pools this bench does not measure, scored as a race between those two alone: on Solana, 77,776 races an hour against roughly 1,200 trades on the bench pool. Feeds that see only our pool were divided by a denominator they could not enter, and read near zero." + - "Race eligibility (since 2026-10-04): only trades our own pool subscription saw are raced. OKX and Birdeye are subscribed per token and see every pool on it, so the book had filled with trades from pools this bench does not measure, scored as a race between those two alone. Measured after the filter: on Solana it discards 576,530 off-pool races an hour and scores 254. Feeds seeing only our pool had been divided by a denominator they could not enter." - "Consequence for the Solana column: `head_lag_seconds` there is the lag behind the first feed in a race, so it moved with the same fix. Figures published before 2026-10-04 for OKX and Birdeye on that column were measured on a wider set of trades than the other four feeds." - "Solana headline: `head_lag_seconds` is the lag behind the first observation of the trade, our node included, floored at 1 ms. GeckoTerminal polls and lands 15 to 100 s later; it is measured against the same first arrival." - "Regions: us-east, eu-west, sgp. Cross-region median reported in the headline." diff --git a/harnesses/aggregator-head-lag/cmd/script/metrics.go b/harnesses/aggregator-head-lag/cmd/script/metrics.go index 25afea98c..2ff6cc59f 100644 --- a/harnesses/aggregator-head-lag/cmd/script/metrics.go +++ b/harnesses/aggregator-head-lag/cmd/script/metrics.go @@ -40,6 +40,9 @@ var ( headLagRaces *prometheus.CounterVec headLagFirstShare *prometheus.GaugeVec headLagRefMatches *prometheus.CounterVec + headLagPoolSeen *prometheus.CounterVec + headLagPoolEntries *prometheus.GaugeVec + headLagRacesVoided *prometheus.CounterVec refClockEntries prometheus.Gauge // Fast-trade latency (for comparison with Pulse V2) @@ -228,6 +231,40 @@ func init() { ) prometheus.MustRegister(headLagRefMatches) + // What the pool subscription actually sees, per chain. The race filter + // reads this set, and until now nothing published its rate, so a filter + // that passed everything and a pool that is genuinely busy produced the + // same observable. + headLagPoolSeen = prometheus.NewCounterVec( + prometheus.CounterOpts{ + Name: "head_lag_pool_observed_total", + Help: "Trades observed by the bench pool subscription, per chain.", + }, + []string{"chain"}, + ) + prometheus.MustRegister(headLagPoolSeen) + + headLagPoolEntries = prometheus.NewGaugeVec( + prometheus.GaugeOpts{ + Name: "head_lag_pool_trades_entries", + Help: "Live entries in the pool membership set (TTL bounded).", + }, + []string{}, + ) + prometheus.MustRegister(headLagPoolEntries) + + // Why a race was discarded. "off_pool" is the filter doing its job; + // "single_participant" is the pre-existing rule. A population that stays + // high with off_pool near zero means the filter is not biting. + headLagRacesVoided = prometheus.NewCounterVec( + prometheus.CounterOpts{ + Name: "head_lag_races_voided_total", + Help: "Races closed without being scored, by reason.", + }, + []string{"chain", "reason"}, + ) + prometheus.MustRegister(headLagRacesVoided) + // Per-transaction race (race.go): who reported the trade first, and // how often, our node reference included as aggregator="reference". headLagFirst = prometheus.NewCounterVec( @@ -516,6 +553,18 @@ func RecordHeadLagRefMiss(aggregator, chain, region string) { } // RecordRefClockSize publishes the reference window occupancy. +func RecordPoolObserved(chain string) { + headLagPoolSeen.WithLabelValues(chain).Inc() +} + +func RecordPoolSetSize(n int) { + headLagPoolEntries.WithLabelValues().Set(float64(n)) +} + +func RecordRaceVoided(chain, reason string) { + headLagRacesVoided.WithLabelValues(chain, reason).Inc() +} + func RecordRefClockSize(n int) { refClockEntries.Set(float64(n)) } // RecordBlockchainHead records the current blockchain head block number diff --git a/harnesses/aggregator-head-lag/cmd/script/race.go b/harnesses/aggregator-head-lag/cmd/script/race.go index 9218655ff..80efdb33e 100644 --- a/harnesses/aggregator-head-lag/cmd/script/race.go +++ b/harnesses/aggregator-head-lag/cmd/script/race.go @@ -195,6 +195,7 @@ func (b *raceBook) closeLocked(e *raceEntry, k string) { e.closed = true e.void = true e.closedAt = time.Now() + RecordRaceVoided(e.chain, "off_pool") return } // Our own node takes part when it saw the trade. It also validates @@ -214,6 +215,7 @@ func (b *raceBook) closeLocked(e *raceEntry, k string) { e.closed = true e.void = true e.closedAt = time.Now() + RecordRaceVoided(e.chain, "single_participant") return } providers := 0 @@ -240,18 +242,18 @@ func (b *raceBook) closeLocked(e *raceEntry, k string) { } } winners := []string{} + referenceFirst := false for _, o := range e.obs { delta := o.at.Sub(e.t0) if o.aggregator == "reference" { if !providerT0.IsZero() && providerT0.Sub(o.at) > raceTie { - headLagFirst.WithLabelValues("reference", e.chain, e.region).Inc() + referenceFirst = true } continue } b.note(e.chain, e.region, o.aggregator) if o.at.Sub(providerT0) <= raceTie { winners = append(winners, o.aggregator) - headLagFirst.WithLabelValues(o.aggregator, e.chain, e.region).Inc() } if raceChains[e.chain] { RecordHeadLag(o.aggregator, e.chain, o.lagBlocks, raceLagSeconds(delta), e.region, hash) @@ -262,9 +264,21 @@ func (b *raceBook) closeLocked(e *raceEntry, k string) { if providers < 2 { // One feed against our node only: lags are recorded above, but // there was no race between feeds to score. + RecordRaceVoided(e.chain, "single_participant") return } headLagRaces.WithLabelValues(e.chain, e.region).Inc() + // Awarded only now, past the single-participant return. These used to be + // incremented inside the loop above, so a trade one feed alone reported + // handed it a first that never entered the race count or the share + // history: a numerator moving without its denominator, which is the same + // defect that put OKX at 95.93% on Solana. + if referenceFirst { + headLagFirst.WithLabelValues("reference", e.chain, e.region).Inc() + } + for _, w := range winners { + headLagFirst.WithLabelValues(w, e.chain, e.region).Inc() + } // Rolling 24 h share. pk := e.chain + "|" + e.region diff --git a/harnesses/aggregator-head-lag/cmd/script/race_test.go b/harnesses/aggregator-head-lag/cmd/script/race_test.go index c3c934845..29f4b022d 100644 --- a/harnesses/aggregator-head-lag/cmd/script/race_test.go +++ b/harnesses/aggregator-head-lag/cmd/script/race_test.go @@ -168,3 +168,38 @@ func TestTokenScopedProvidersAreTheOnesNeedingPoolFiltering(t *testing.T) { } } } + +// A trade one feed alone reported must not award it a first. +// +// The counter used to be incremented inside the winners loop, which runs +// before the single-participant return, so such a trade handed that feed a +// first while never entering the race count or the share history. A numerator +// moving without its denominator is the same defect that published OKX at +// 95.93% on Solana, and it survived the first fix. +func TestSingleProviderRaceAwardsNoFirst(t *testing.T) { + withPoolTrades(t, func(p *refClock) { + p.observe("solana", "LONEHASH", time.Now()) + }) + // A reference observation makes len(obs) == 2, so the race passes the + // participant guard and reaches the providers < 2 return. That is exactly + // the path that used to award the free first. + savedRef := reference + reference = &refClock{seen: map[string]refEntry{}} + reference.observe("solana", "LONEHASH", time.Now()) + t.Cleanup(func() { reference = savedRef }) + + b := newTestBook() + now := time.Now() + b.observe("okx", "solana", "eu-west", "LONEHASH", now, 0) + + b.resolve(now.Add(raceWindow + time.Second)) + + e := b.entries[raceKey("solana", "eu-west", "LONEHASH")] + if e == nil || !e.closed { + t.Fatal("race did not close") + } + if len(b.history["solana|eu-west"]) != 0 { + t.Error("a single-provider race entered the share history, so it would " + + "sit in every other feed's denominator") + } +} diff --git a/harnesses/aggregator-head-lag/cmd/script/reference_monitor.go b/harnesses/aggregator-head-lag/cmd/script/reference_monitor.go index a293129bc..a1d2b051a 100644 --- a/harnesses/aggregator-head-lag/cmd/script/reference_monitor.go +++ b/harnesses/aggregator-head-lag/cmd/script/reference_monitor.go @@ -178,6 +178,7 @@ func runReferenceMonitor(stopChan <-chan struct{}) { reference.sweep() poolTrades.sweep() RecordRefClockSize(reference.size()) + RecordPoolSetSize(poolTrades.size()) } } }() @@ -304,6 +305,7 @@ func refConnect(p HeadLagPool, url string, stopChan <-chan struct{}) error { } reference.observe(p.ChainName, r.Value.Signature, now) poolTrades.observe(p.ChainName, r.Value.Signature, now) + RecordPoolObserved(p.ChainName) continue } @@ -320,6 +322,7 @@ func refConnect(p HeadLagPool, url string, stopChan <-chan struct{}) error { } reference.observe(p.ChainName, r.TransactionHash, now) poolTrades.observe(p.ChainName, r.TransactionHash, now) + RecordPoolObserved(p.ChainName) } }