Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion benchmarks/aggregator-head-lag.yml
Original file line number Diff line number Diff line change
Expand Up @@ -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."
Expand Down
49 changes: 49 additions & 0 deletions harnesses/aggregator-head-lag/cmd/script/metrics.go
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down Expand Up @@ -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(
Expand Down Expand Up @@ -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
Expand Down
18 changes: 16 additions & 2 deletions harnesses/aggregator-head-lag/cmd/script/race.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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
Expand All @@ -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)
Expand All @@ -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
Expand Down
35 changes: 35 additions & 0 deletions harnesses/aggregator-head-lag/cmd/script/race_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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")
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -178,6 +178,7 @@ func runReferenceMonitor(stopChan <-chan struct{}) {
reference.sweep()
poolTrades.sweep()
RecordRefClockSize(reference.size())
RecordPoolSetSize(poolTrades.size())
}
}
}()
Expand Down Expand Up @@ -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
}

Expand All @@ -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)
}
}

Expand Down
Loading