diff --git a/pkg/atproto/firehose.go b/pkg/atproto/firehose.go index 3477d34ff..b42a8224e 100644 --- a/pkg/atproto/firehose.go +++ b/pkg/atproto/firehose.go @@ -273,32 +273,54 @@ func (atsync *ATProtoSynchronizer) connectRelay(ctx context.Context, relay strin } defer con.Close() - spmetrics.FirehoseRelaysConnected.WithLabelValues(relay).Set(1) - defer spmetrics.FirehoseRelaysConnected.WithLabelValues(relay).Set(0) + protocol := relayProtocol(relay) + spmetrics.FirehoseRelaysConnected.WithLabelValues(relay, protocol).Set(1) + defer spmetrics.FirehoseRelaysConnected.WithLabelValues(relay, protocol).Set(0) streamCtx, cancel := context.WithCancel(ctx) defer cancel() - rsc := atsync.repoStreamCallbacks(ctx, cursor, cancel) + rsc := atsync.repoStreamCallbacks(ctx, relay, cursor, cancel) scheduler := parallel.NewScheduler(10, 100, relay, rsc.EventHandler) log.Log(ctx, "connected to relay firehose") return events.HandleRepoStream(streamCtx, con, scheduler, nil) } +// relayProtocol classifies a relay URL as "moq" (moqt:// and aliases) or +// "websocket" (everything else), for per-protocol metric labels. +func relayProtocol(relay string) string { + if u, err := url.Parse(relay); err == nil { + switch u.Scheme { + case "moqt", "moql", "moq", "moqs": + return "moq" + } + } + return "websocket" +} + // repoStreamCallbacks builds the event callbacks shared by the WebSocket and -// MoQ transports: liveness marking, cursor advance, cross-relay dedup, and -// spawning the indexing handlers on the parent ctx. cancel ends the current -// connection when the relay sends an error frame. The handlers run on ctx (not -// the per-connection context) so an in-flight commit keeps indexing across a -// reconnect. -func (atsync *ATProtoSynchronizer) repoStreamCallbacks(ctx context.Context, cursor *relayCursor, cancel context.CancelFunc) *events.RepoStreamCallbacks { +// MoQ transports: per-relay metrics, liveness marking, cursor advance, +// cross-relay dedup, and spawning the indexing handlers on the parent ctx. +// cancel ends the current connection when the relay sends an error frame. The +// handlers run on ctx (not the per-connection context) so an in-flight commit +// keeps indexing across a reconnect. +func (atsync *ATProtoSynchronizer) repoStreamCallbacks(ctx context.Context, relay string, cursor *relayCursor, cancel context.CancelFunc) *events.RepoStreamCallbacks { + protocol := relayProtocol(relay) + // observeSeq advances the cursor and republishes the per-relay high-water + // seq gauge; both relays carry the upstream's seq, so the gauge's cross-relay + // difference measures how far apart the relays are. + observeSeq := func(seq int64) { + cursor.observe(seq) + spmetrics.FirehoseRelayHighSeq.WithLabelValues(relay, protocol).Set(float64(cursor.highSeq())) + } return &events.RepoStreamCallbacks{ RepoCommit: func(evt *comatproto.SyncSubscribeRepos_Commit) error { atsync.markSeen() - cursor.observe(evt.Seq) + observeSeq(evt.Seq) + spmetrics.FirehoseEventsReceivedTotal.WithLabelValues(relay, protocol, "commit").Inc() if atsync.commitDedup.seen(evt.Commit.String()) { - spmetrics.FirehoseEventsDedupedTotal.WithLabelValues("commit").Inc() + spmetrics.FirehoseEventsDedupedTotal.WithLabelValues(relay, protocol, "commit").Inc() return nil } go atsync.handleCommitEventOps(ctx, evt) @@ -306,9 +328,10 @@ func (atsync *ATProtoSynchronizer) repoStreamCallbacks(ctx context.Context, curs }, RepoIdentity: func(evt *comatproto.SyncSubscribeRepos_Identity) error { atsync.markSeen() - cursor.observe(evt.Seq) + observeSeq(evt.Seq) + spmetrics.FirehoseEventsReceivedTotal.WithLabelValues(relay, protocol, "identity").Inc() if atsync.identityDedup.seen(identityDedupKey(evt)) { - spmetrics.FirehoseEventsDedupedTotal.WithLabelValues("identity").Inc() + spmetrics.FirehoseEventsDedupedTotal.WithLabelValues(relay, protocol, "identity").Inc() return nil } go atsync.handleIdentityEventOps(ctx, evt) diff --git a/pkg/atproto/firehose_cursor.go b/pkg/atproto/firehose_cursor.go index 2bfcb9f9c..2cd013ffd 100644 --- a/pkg/atproto/firehose_cursor.go +++ b/pkg/atproto/firehose_cursor.go @@ -86,6 +86,9 @@ func (rc *relayCursor) param() (int64, bool) { return v, v > 0 } +// highSeq returns the high-water upstream sequence number observed so far. +func (rc *relayCursor) highSeq() int64 { return rc.latest.Load() } + // observeGroup advances the high-water MoQ group sequence. Called on every // frame received from a moqt:// relay (concurrency-safe). func (rc *relayCursor) observeGroup(seq uint64) { diff --git a/pkg/atproto/firehose_moq.go b/pkg/atproto/firehose_moq.go index b11b79128..7b42bc133 100644 --- a/pkg/atproto/firehose_moq.go +++ b/pkg/atproto/firehose_moq.go @@ -49,10 +49,11 @@ func (atsync *ATProtoSynchronizer) connectRelayMoq(ctx context.Context, relay st } defer sub.Close() - spmetrics.FirehoseRelaysConnected.WithLabelValues(relay).Set(1) - defer spmetrics.FirehoseRelaysConnected.WithLabelValues(relay).Set(0) + protocol := relayProtocol(relay) + spmetrics.FirehoseRelaysConnected.WithLabelValues(relay, protocol).Set(1) + defer spmetrics.FirehoseRelaysConnected.WithLabelValues(relay, protocol).Set(0) - rsc := atsync.repoStreamCallbacks(ctx, cursor, cancel) + rsc := atsync.repoStreamCallbacks(ctx, relay, cursor, cancel) scheduler := parallel.NewScheduler(10, 100, relay, rsc.EventHandler) defer scheduler.Shutdown() diff --git a/pkg/atproto/firehose_multirelay_test.go b/pkg/atproto/firehose_multirelay_test.go index d402ff024..5a28e9957 100644 --- a/pkg/atproto/firehose_multirelay_test.go +++ b/pkg/atproto/firehose_multirelay_test.go @@ -11,13 +11,11 @@ import ( lexutil "github.com/bluesky-social/indigo/lex/util" "github.com/bluesky-social/indigo/util" "github.com/prometheus/client_golang/prometheus" - dto "github.com/prometheus/client_model/go" "github.com/stretchr/testify/require" "stream.place/streamplace/pkg/bus" "stream.place/streamplace/pkg/config" "stream.place/streamplace/pkg/devenv" "stream.place/streamplace/pkg/model" - "stream.place/streamplace/pkg/spmetrics" "stream.place/streamplace/pkg/statedb" "stream.place/streamplace/pkg/streamplace" ) @@ -64,7 +62,7 @@ func TestMultiRelayDedup(t *testing.T) { // Sanity-check the fan-out: two relays, no self. require.Len(t, atsync.relayHosts(), 2) - dedupedBefore := counterValue(t, spmetrics.FirehoseEventsDedupedTotal.WithLabelValues("commit")) + dedupedBefore := dedupedCommitTotal(t) go func() { _ = atsync.StartFirehose(ctx) @@ -99,7 +97,7 @@ func TestMultiRelayDedup(t *testing.T) { // And the duplicates the second relay delivered were dropped by the deduper. err = untilNoErrors(t, func() error { - deduped := counterValue(t, spmetrics.FirehoseEventsDedupedTotal.WithLabelValues("commit")) + deduped := dedupedCommitTotal(t) if deduped <= dedupedBefore { return fmt.Errorf("expected commit dedup counter to climb (before=%v now=%v)", dedupedBefore, deduped) } @@ -108,9 +106,24 @@ func TestMultiRelayDedup(t *testing.T) { require.NoError(t, err) } -func counterValue(t *testing.T, c prometheus.Counter) float64 { +// dedupedCommitTotal sums the commit dedup counter across every relay/protocol +// series (a duplicate is attributed to whichever relay delivered it second). +func dedupedCommitTotal(t *testing.T) float64 { t.Helper() - var m dto.Metric - require.NoError(t, c.Write(&m)) - return m.GetCounter().GetValue() + mfs, err := prometheus.DefaultGatherer.Gather() + require.NoError(t, err) + var sum float64 + for _, mf := range mfs { + if mf.GetName() != "streamplace_firehose_events_deduped_total" { + continue + } + for _, m := range mf.GetMetric() { + for _, l := range m.GetLabel() { + if l.GetName() == "kind" && l.GetValue() == "commit" { + sum += m.GetCounter().GetValue() + } + } + } + } + return sum } diff --git a/pkg/spmetrics/spmetrics.go b/pkg/spmetrics/spmetrics.go index b62673baa..af784f623 100644 --- a/pkg/spmetrics/spmetrics.go +++ b/pkg/spmetrics/spmetrics.go @@ -103,23 +103,45 @@ var LabelerFirehosesConnected = promauto.NewGaugeVec(prometheus.GaugeOpts{ Help: "number of currently connected labeler firehoses", }, []string{"labeler"}) -// FirehoseRelaysConnected is 1 while a relay's subscribeRepos websocket is -// connected and 0 while it is reconnecting, labeled by relay host. With -// multi-relay support this shows at a glance how many of the configured -// relays are currently feeding us. +// FirehoseRelaysConnected is 1 while a relay's firehose is connected and 0 +// while it is reconnecting, labeled by relay host and protocol ("websocket" or +// "moq"). With multi-relay support this shows at a glance how many of the +// configured relays are currently feeding us. var FirehoseRelaysConnected = promauto.NewGaugeVec(prometheus.GaugeOpts{ Name: "streamplace_firehose_relays_connected", - Help: "1 if the relay's firehose websocket is currently connected, else 0", -}, []string{"relay"}) - -// FirehoseEventsDedupedTotal counts events dropped because the same commit -// (or identity update) already arrived from another relay. Labeled by event -// kind ("commit" / "identity"). A high count is expected and healthy — it is -// the redundant traffic we are paying for resilience. + Help: "1 if the relay's firehose is currently connected, else 0", +}, []string{"relay", "protocol"}) + +// FirehoseEventsReceivedTotal counts every event a relay delivered, before +// cross-relay dedup, labeled by relay, protocol ("websocket"/"moq"), and kind. +// Comparing rates across relays tells you whether they carry the same volume; +// received minus deduped is each relay's first-seen (unique) contribution. +var FirehoseEventsReceivedTotal = promauto.NewCounterVec(prometheus.CounterOpts{ + Name: "streamplace_firehose_events_received_total", + Help: "firehose events delivered by each relay before dedup, by relay/protocol/kind", +}, []string{"relay", "protocol", "kind"}) + +// FirehoseEventsDedupedTotal counts events dropped because the same commit (or +// identity update) already arrived from another relay, labeled by the relay that +// delivered the duplicate, its protocol, and event kind. A relay whose deduped +// count tracks its received count is fully redundant (a healthy mirror); a relay +// with a high first-seen count (received minus deduped) is contributing events +// the others are not — i.e. the relays are diverging. var FirehoseEventsDedupedTotal = promauto.NewCounterVec(prometheus.CounterOpts{ Name: "streamplace_firehose_events_deduped_total", - Help: "firehose events dropped as cross-relay duplicates, by kind", -}, []string{"kind"}) + Help: "firehose events dropped as cross-relay duplicates, by relay/protocol/kind", +}, []string{"relay", "protocol", "kind"}) + +// FirehoseRelayHighSeq is the highest upstream sequence number seen from each +// relay. A WebSocket bsky relay and a MoQ relay that bridges the same upstream +// both carry bsky's original seq, so the difference between two relays' high_seq +// is a direct, units-of-events measure of how far apart they are: a small, +// roughly-constant skew means they are in sync (just a latency offset); a +// growing gap means one relay is falling behind or stuck. +var FirehoseRelayHighSeq = promauto.NewGaugeVec(prometheus.GaugeOpts{ + Name: "streamplace_firehose_relay_high_seq", + Help: "highest upstream sequence number seen from each relay (cross-relay diff = parity/lag)", +}, []string{"relay", "protocol"}) // --- isolated ingest workers ------------------------------------------------