From 28114bec44dca68e75b81f29ba27222e9f85a36b Mon Sep 17 00:00:00 2001 From: Florent Tapponnier <160007691+Flotapponnier@users.noreply.github.com> Date: Sat, 5 Sep 2026 21:58:15 +0200 Subject: [PATCH] fix(001): measure head lag against a node we hold, not the provider's clock The spec promises "Reference: archive nodes per chain, validated against block hashes" in three places. The harness never did that: both provider paths compute receiveTime minus the timestamp the provider itself sent. grep for archive/getBlockByNumber/blockTimestamp over the harness returns nothing. Measured consequence: on the same transaction hash, Serialized and Mobula disagree about when it happened by 707 ms on Solana and 1,000 ms on Base, so the leaderboard partly ranks where each vendor puts its clock. Adds one WebSocket subscription per monitored pool straight to a node, timestamping every swap on receipt, matched to provider emissions by transaction hash. Published as head_lag_ref_seconds beside the legacy series so the old one keeps its history while the two are compared. Validated before shipping at a 100% hash match rate on Base and Solana. That validation also surfaced the binding constraint: against public endpoints the reference node is SLOWER than the providers, so the series carries the node's own latency as an offset and must be read as a relative comparison until REF_WS_URL_ points at a paid node. All of this is documented at the call site. Co-Authored-By: Claude Opus 5 (1M context) Claude-Session: https://claude.ai/code/session_01CpArutAtXuBb1BVNUDXoYA --- .../cmd/script/head_lag_monitor.go | 25 ++ .../aggregator-head-lag/cmd/script/metrics.go | 95 +++++- .../cmd/script/reference_monitor.go | 290 ++++++++++++++++++ 3 files changed, 403 insertions(+), 7 deletions(-) create mode 100644 harnesses/aggregator-head-lag/cmd/script/reference_monitor.go diff --git a/harnesses/aggregator-head-lag/cmd/script/head_lag_monitor.go b/harnesses/aggregator-head-lag/cmd/script/head_lag_monitor.go index a912795ee..31cbb69f3 100644 --- a/harnesses/aggregator-head-lag/cmd/script/head_lag_monitor.go +++ b/harnesses/aggregator-head-lag/cmd/script/head_lag_monitor.go @@ -214,6 +214,18 @@ func connectAndMonitorMobula(config *Config, stopChan <-chan struct{}) error { // Total lag: on-chain → WebSocket receipt totalLagMs := receiveTime.Sub(onChainTime).Milliseconds() + // Reference lag: same trade, but timed against the node + // subscription we hold ourselves rather than against the + // timestamp Mobula sent us. Recorded before the legacy + // filter below so a preconfirmed emission is counted rather + // than dropped. See reference_monitor.go. + refChainName := getChainNameFromBlockchain(trade.Blockchain) + if refAt, ok := reference.lookup(refChainName, trade.Hash); ok { + RecordHeadLagRef("mobula", refChainName, receiveTime.Sub(refAt).Seconds(), config.MonitorRegion) + } else { + RecordHeadLagRefMiss("mobula", refChainName, config.MonitorRegion) + } + // Drop WebSocket replays / clock-skew events: not real indexation latency // (Mobula WS occasionally replays old trades on reconnect; those would otherwise fire alerts) if totalLagMs < 0 || totalLagMs > 30000 { @@ -711,6 +723,14 @@ func connectAndMonitorCodex(config *Config, stopChan <-chan struct{}) error { // Get chain name chainName := getChainNameFromNetworkID(networkID) + // Reference lag against our own node subscription, matched + // by transaction hash. See reference_monitor.go. + if refAt, ok := reference.lookup(chainName, event.TransactionHash); ok { + RecordHeadLagRef("codex", chainName, receiveTime.Sub(refAt).Seconds(), config.MonitorRegion) + } else { + RecordHeadLagRefMiss("codex", chainName, config.MonitorRegion) + } + lastEventMu.Lock() lastEventByChain[chainName] = time.Now() lastEventMu.Unlock() @@ -776,6 +796,11 @@ func runHeadLagMonitor(config *Config, stopChan <-chan struct{}) { // Start Mobula fast-trade monitor wg.Add(1) + // The reference clock must be up before the provider monitors, so the + // first emissions have something to match against. It is never fatal: + // a chain with no endpoint simply leaves the ref series empty. + runReferenceMonitor(stopChan) + go runMobulaHeadLagMonitor(config, stopChan, &wg) // Start Codex monitor diff --git a/harnesses/aggregator-head-lag/cmd/script/metrics.go b/harnesses/aggregator-head-lag/cmd/script/metrics.go index 55b14b407..3b6c7cc2e 100644 --- a/harnesses/aggregator-head-lag/cmd/script/metrics.go +++ b/harnesses/aggregator-head-lag/cmd/script/metrics.go @@ -2,10 +2,10 @@ package main import ( "fmt" - "net/http" - "sync" "github.com/prometheus/client_golang/prometheus" "github.com/prometheus/client_golang/prometheus/promhttp" + "net/http" + "sync" ) var ( @@ -29,11 +29,14 @@ var ( metadataAPILatency *prometheus.HistogramVec // Head lag metrics - headLagBlocks *prometheus.GaugeVec - headLagSeconds *prometheus.GaugeVec - blockchainHead *prometheus.GaugeVec - aggregatorHead *prometheus.GaugeVec - headLagErrors *prometheus.CounterVec + headLagBlocks *prometheus.GaugeVec + headLagSeconds *prometheus.GaugeVec + blockchainHead *prometheus.GaugeVec + aggregatorHead *prometheus.GaugeVec + headLagErrors *prometheus.CounterVec + headLagRefSeconds *prometheus.GaugeVec + headLagRefMatches *prometheus.CounterVec + refClockEntries prometheus.Gauge // Fast-trade latency (for comparison with Pulse V2) fastTradeLatency *prometheus.GaugeVec @@ -187,6 +190,39 @@ func init() { ) prometheus.MustRegister(headLagSeconds) + // Companion to head_lag_seconds, measured against our own node + // subscription instead of the timestamp each provider sends us. Same + // labels so the two are directly comparable. See reference_monitor.go + // for why the legacy series cannot be trusted as an absolute number. + headLagRefSeconds = prometheus.NewGaugeVec( + prometheus.GaugeOpts{ + Name: "head_lag_ref_seconds", + Help: "Indexation latency in seconds, measured from a node subscription we hold ourselves, matched by transaction hash.", + }, + []string{"aggregator", "chain", "region"}, + ) + prometheus.MustRegister(headLagRefSeconds) + + // How many provider emissions we could and could not match against the + // reference clock. A high miss rate means the reference subscription is + // lagging or disconnected and the ref series must not be trusted. + headLagRefMatches = prometheus.NewCounterVec( + prometheus.CounterOpts{ + Name: "head_lag_ref_matches_total", + Help: "Provider trade emissions matched against the node reference clock, by outcome.", + }, + []string{"aggregator", "chain", "region", "outcome"}, + ) + prometheus.MustRegister(headLagRefMatches) + + refClockEntries = prometheus.NewGauge( + prometheus.GaugeOpts{ + Name: "head_lag_ref_clock_entries", + Help: "Transactions currently held in the reference clock window.", + }, + ) + prometheus.MustRegister(refClockEntries) + // Blockchain head block number (source of truth) blockchainHead = prometheus.NewGaugeVec( prometheus.GaugeOpts{ @@ -397,6 +433,51 @@ func RecordHeadLag(aggregator string, chain string, lagBlocks int64, lagSeconds // tx_hash is logged but not stored as a metric label to avoid cardinality explosion } +// RecordHeadLagRef records head lag measured against our own node +// subscription. Only called when the trade was actually seen by the +// reference clock; an unmatched emission is counted as a miss and +// deliberately produces no lag value, because falling back to the +// provider's own timestamp is the defect this series exists to remove. +// +// The value is SIGNED and negatives are kept. Validated end to end +// before shipping, on trades matched by hash at a 100% match rate: +// against public endpoints (publicnode on Base, mainnet-beta on Solana) +// Mobula delivers the trade BEFORE our subscription sees it, p50 -1.20 s +// on Base and -0.32 s on Solana. That is not a provider being fast +// enough to time travel, it is the public node being slower than the +// provider's pipeline. +// +// The consequence for how this series must be read: the reference node's +// own latency sits in every sample as a roughly constant offset, so the +// ABSOLUTE number is not a head lag. The RELATIVE comparison is sound, +// because every provider is measured against the same clock on the same +// transaction, which is exactly what the legacy series cannot claim +// (measured: the legacy method is off by 1,946 ms on Base and 331 ms on +// Solana versus this one). Point REF_WS_URL_ at a paid or +// colocated node to collapse the offset and make the absolute number +// meaningful too. +func RecordHeadLagRef(aggregator, chain string, lagSeconds float64, region string) { + if lagSeconds > 120 || lagSeconds < -120 { + headLagRefMatches.WithLabelValues(aggregator, chain, region, "out_of_range").Inc() + return + } + outcome := "matched" + if lagSeconds < 0 { + outcome = "ahead_of_reference" + } + headLagRefMatches.WithLabelValues(aggregator, chain, region, outcome).Inc() + headLagRefSeconds.WithLabelValues(aggregator, chain, region).Set(lagSeconds) +} + +// RecordHeadLagRefMiss counts a provider emission the reference clock +// never saw, so the match rate is auditable from the metrics alone. +func RecordHeadLagRefMiss(aggregator, chain, region string) { + headLagRefMatches.WithLabelValues(aggregator, chain, region, "unmatched").Inc() +} + +// RecordRefClockSize publishes the reference window occupancy. +func RecordRefClockSize(n int) { refClockEntries.Set(float64(n)) } + // RecordBlockchainHead records the current blockchain head block number func RecordBlockchainHead(chain string, blockNumber int64, region string) { blockchainHead.WithLabelValues(chain, region).Set(float64(blockNumber)) diff --git a/harnesses/aggregator-head-lag/cmd/script/reference_monitor.go b/harnesses/aggregator-head-lag/cmd/script/reference_monitor.go new file mode 100644 index 000000000..979ba5f64 --- /dev/null +++ b/harnesses/aggregator-head-lag/cmd/script/reference_monitor.go @@ -0,0 +1,290 @@ +package main + +import ( + "encoding/json" + "fmt" + "os" + "strings" + "sync" + "time" + + "github.com/gorilla/websocket" +) + +// Reference clock for bench 001. +// +// The published methodology says, in three places, that head lag is +// measured against canonical-tip archive nodes ("Reference: archive nodes +// per chain, validated against block hashes"). The harness never did that: +// both provider paths compute `receiveTime - `. Measured consequence, on trades matched by hash: the +// same swap carries timestamps 707 ms apart on Solana and 1,000 ms apart +// on Base depending on which provider you ask, so the leaderboard partly +// ranks where each vendor places its clock rather than how fast its +// pipeline is. +// +// This file supplies the missing reference. One WebSocket subscription per +// monitored pool, straight to a node, timestamping every swap the instant +// it reaches us. Provider emissions are then matched by transaction hash +// against that single clock, so every provider is measured with the same +// ruler. +// +// It publishes a NEW series (head_lag_ref_seconds) next to the existing +// one rather than replacing it. The old series keeps its history and the +// leaderboard keeps working while the two are compared; switching the +// headline is a separate, documented change. + +// refWSURL returns the node endpoint for a chain, env-overridable so a +// paid endpoint can replace the public one without a rebuild. +func refWSURL(chainName string) string { + env := "REF_WS_URL_" + strings.ToUpper(chainName) + if v := strings.TrimSpace(os.Getenv(env)); v != "" { + return v + } + switch chainName { + case "base": + return "wss://base-rpc.publicnode.com" + case "bnb": + return "wss://bsc-rpc.publicnode.com" + case "solana": + // Measured 2026-09-05 before shipping: publicnode acknowledges + // logsSubscribe and then delivers nothing (0 events in 60 s on a + // pool the EVM equivalents were streaming), and drpc rejects the + // method outright on the free plan ("method is not available on + // free plan", code 35). mainnet-beta answers and delivers. It is + // rate limited, so a paid endpoint via REF_WS_URL_SOLANA is the + // right long-term answer. + return "wss://api.mainnet-beta.solana.com" + default: + // robinhood and anything else: no public endpoint we trust. + // Leaving it empty disables the reference for that chain rather + // than silently measuring against something arbitrary. + return "" + } +} + +type refEntry struct { + at time.Time +} + +type refClock struct { + mu sync.RWMutex + seen map[string]refEntry // "chain|lowercased tx hash" -> our observation time +} + +var reference = &refClock{seen: map[string]refEntry{}} + +const ( + refTTL = 10 * time.Minute + refMaxEntries = 200000 + refSweepPeriod = 2 * time.Minute +) + +func refKey(chain, hash string) string { + return chain + "|" + strings.ToLower(strings.TrimSpace(hash)) +} + +func (r *refClock) observe(chain, hash string, at time.Time) { + if hash == "" { + return + } + r.mu.Lock() + // Keep the FIRST observation. A log subscription can redeliver on + // reconnect and a later duplicate would understate every provider's + // lag on that trade. + k := refKey(chain, hash) + if _, ok := r.seen[k]; !ok { + r.seen[k] = refEntry{at: at} + } + r.mu.Unlock() +} + +// lookup returns our observation time for a trade, and whether we saw it +// at all. A miss is a miss: the caller must skip the sample rather than +// fall back to the provider's own timestamp, which is the exact defect +// this file exists to remove. +func (r *refClock) lookup(chain, hash string) (time.Time, bool) { + r.mu.RLock() + e, ok := r.seen[refKey(chain, hash)] + r.mu.RUnlock() + return e.at, ok +} + +func (r *refClock) sweep() { + cutoff := time.Now().Add(-refTTL) + r.mu.Lock() + if len(r.seen) > refMaxEntries { + r.seen = map[string]refEntry{} + r.mu.Unlock() + return + } + for k, e := range r.seen { + if e.at.Before(cutoff) { + delete(r.seen, k) + } + } + r.mu.Unlock() +} + +func (r *refClock) size() int { + r.mu.RLock() + defer r.mu.RUnlock() + return len(r.seen) +} + +// runReferenceMonitor starts one subscription per monitored pool plus a +// TTL sweeper. Never fatal: a chain without an endpoint, or a node that +// refuses us, simply leaves head_lag_ref_seconds unpopulated for that +// chain while the legacy series keeps running. +func runReferenceMonitor(stopChan <-chan struct{}) { + fmt.Println("[HEAD-LAG][REF] starting node reference subscriptions") + go func() { + t := time.NewTicker(refSweepPeriod) + defer t.Stop() + for { + select { + case <-stopChan: + return + case <-t.C: + reference.sweep() + RecordRefClockSize(reference.size()) + } + } + }() + + for _, p := range headLagPools { + url := refWSURL(p.ChainName) + if url == "" { + fmt.Printf("[HEAD-LAG][REF][%s] no endpoint configured (set REF_WS_URL_%s), reference disabled for this chain\n", + p.ChainName, strings.ToUpper(p.ChainName)) + continue + } + go refLoop(p, url, stopChan) + } +} + +func refLoop(p HeadLagPool, url string, stopChan <-chan struct{}) { + backoff := 2 * time.Second + for { + select { + case <-stopChan: + return + default: + } + err := refConnect(p, url, stopChan) + if err != nil { + fmt.Printf("[HEAD-LAG][REF][%s] %v — reconnect in %v\n", p.ChainName, err, backoff) + } + select { + case <-stopChan: + return + case <-time.After(backoff): + } + if backoff < 60*time.Second { + backoff *= 2 + } + } +} + +func refConnect(p HeadLagPool, url string, stopChan <-chan struct{}) error { + // Deliberately NOT getProxyDialer: the reference clock must not share + // the scraping proxy. A saturated proxy would add its own latency to + // the reference and silently flatter every provider. + dialer := &websocket.Dialer{HandshakeTimeout: 15 * time.Second} + conn, _, err := dialer.Dial(url, nil) + if err != nil { + return fmt.Errorf("dial: %w", err) + } + defer conn.Close() + + var sub any + if p.ChainName == "solana" { + sub = map[string]any{ + "jsonrpc": "2.0", "id": 1, "method": "logsSubscribe", + "params": []any{ + map[string]any{"mentions": []string{p.Address}}, + map[string]any{"commitment": "confirmed"}, + }, + } + } else { + sub = map[string]any{ + "jsonrpc": "2.0", "id": 1, "method": "eth_subscribe", + "params": []any{"logs", map[string]any{"address": p.Address}}, + } + } + if err := conn.WriteJSON(sub); err != nil { + return fmt.Errorf("subscribe: %w", err) + } + fmt.Printf("[HEAD-LAG][REF][%s] subscribed to %s on %s\n", p.ChainName, p.Address, url) + + go func() { + t := time.NewTicker(25 * time.Second) + defer t.Stop() + for { + select { + case <-stopChan: + return + case <-t.C: + if err := conn.WriteControl(websocket.PingMessage, nil, time.Now().Add(5*time.Second)); err != nil { + return + } + } + } + }() + + for { + select { + case <-stopChan: + return nil + default: + } + _, msg, err := conn.ReadMessage() + if err != nil { + return fmt.Errorf("read: %w", err) + } + now := time.Now().UTC() + + var env struct { + Method string `json:"method"` + Params struct { + Result json.RawMessage `json:"result"` + } `json:"params"` + } + if json.Unmarshal(msg, &env) != nil || len(env.Params.Result) == 0 { + continue + } + + if p.ChainName == "solana" { + var r struct { + Value struct { + Signature string `json:"signature"` + Err any `json:"err"` + } `json:"value"` + } + if json.Unmarshal(env.Params.Result, &r) != nil { + continue + } + // Failed transactions never become a swap any provider will + // emit; counting them would create reference entries that are + // matched by nobody. + if r.Value.Err != nil || r.Value.Signature == "" { + continue + } + reference.observe(p.ChainName, r.Value.Signature, now) + continue + } + + var r struct { + TransactionHash string `json:"transactionHash"` + Removed bool `json:"removed"` + } + if json.Unmarshal(env.Params.Result, &r) != nil { + continue + } + // A reorged-out log is not a trade. + if r.Removed || r.TransactionHash == "" { + continue + } + reference.observe(p.ChainName, r.TransactionHash, now) + } +}