From 753ef82fdb5c56adfb106c238e0d2eb029242436 Mon Sep 17 00:00:00 2001 From: Seongmin Lee Date: Tue, 14 Jul 2026 23:43:57 +0900 Subject: [PATCH] wip: knotmirror/relay 1:1 mapping with indigo/cmd/relay/relay Signed-off-by: Seongmin Lee --- knotmirror/knotmirror.go | 8 ++ knotmirror/relay/broadcast.go | 221 ++++++++++++++++++++++++++++++++++ knotmirror/relay/ingest.go | 61 ++++++++++ knotmirror/relay/metrics.go | 49 ++++++++ knotmirror/relay/relay.go | 88 ++++++++++++++ 5 files changed, 427 insertions(+) create mode 100644 knotmirror/relay/broadcast.go create mode 100644 knotmirror/relay/ingest.go create mode 100644 knotmirror/relay/metrics.go create mode 100644 knotmirror/relay/relay.go diff --git a/knotmirror/knotmirror.go b/knotmirror/knotmirror.go index c07911d7..485c922c 100644 --- a/knotmirror/knotmirror.go +++ b/knotmirror/knotmirror.go @@ -15,6 +15,7 @@ import ( "tangled.org/core/knotmirror/db" "tangled.org/core/knotmirror/knotstream" "tangled.org/core/knotmirror/models" + "tangled.org/core/knotmirror/relay" "tangled.org/core/knotmirror/repoindexer" "tangled.org/core/knotmirror/stream/eventmgr" "tangled.org/core/knotmirror/xrpc" @@ -40,6 +41,7 @@ func Run(ctx context.Context, cfg *config.Config) error { logger.Error("failed to create redis resolver for admin, falling back to default", "err", err) resolver = idresolver.DefaultResolver(cfg.PlcUrl) } + dir := resolver.Directory() // NOTE: using plain git-cli for clone/fetch as go-git is too memory-intensive. gitm := NewCliGitMirrorManager(cfg.GitRepoBasePath, cfg.KnotUseSSL) @@ -64,6 +66,12 @@ func Run(ctx context.Context, cfg *config.Config) error { evtman := eventmgr.NewEventManager(persister) + logger.Info("constructing relay service") + r, err := relay.NewRelay(db, evtman, dir, relayConfig) + if err != nil { + return err + } + knotstream := knotstream.NewKnotStream(logger, db, cfg) crawler := NewCrawler(logger, db) resyncer := NewResyncer(logger, db, gitm, indexScheduler, cfg) diff --git a/knotmirror/relay/broadcast.go b/knotmirror/relay/broadcast.go new file mode 100644 index 00000000..b90f925d --- /dev/null +++ b/knotmirror/relay/broadcast.go @@ -0,0 +1,221 @@ +package relay + +import ( + "context" + "fmt" + "net/http" + "sync" + "time" + + "tangled.org/core/knotmirror/stream" + + "github.com/gorilla/websocket" + promclient "github.com/prometheus/client_golang/prometheus" + dto "github.com/prometheus/client_model/go" +) + +type SocketConsumer struct { + UserAgent string + RemoteAddr string + ConnectedAt time.Time + EventsSent promclient.Counter +} + +func (r *Relay) registerConsumer(c *SocketConsumer) uint64 { + r.consumersLk.Lock() + defer r.consumersLk.Unlock() + + id := r.nextConsumerID + r.nextConsumerID++ + + r.consumers[id] = c + + return id +} + +func (r *Relay) cleanupConsumer(id uint64) { + r.consumersLk.Lock() + defer r.consumersLk.Unlock() + + c := r.consumers[id] + + var m = &dto.Metric{} + if err := c.EventsSent.Write(m); err != nil { + r.Logger.Error("failed to get sent counter", "err", err) + } + + r.Logger.Info("consumer disconnected", + "consumer_id", id, + "remote_addr", c.RemoteAddr, + "user_agent", c.UserAgent, + "events_sent", m.Counter.GetValue()) + + delete(r.consumers, id) +} + +var wsUpgrader = websocket.Upgrader{ + ReadBufferSize: 10_000, + WriteBufferSize: 10_000, +} + +// Main HTTP request handler for clients connecting to the firehose (org.tangled.sync.subscribeRepos) +func (r *Relay) HandleSubscribeRepos(resp http.ResponseWriter, req *http.Request, since *int64, realIP string) error { + + ctx, cancel := context.WithCancel(req.Context()) + defer cancel() + + conn, err := wsUpgrader.Upgrade(resp, req, resp.Header()) + if err != nil { + return fmt.Errorf("upgrading websocket: %w", err) + } + + defer func() { + _ = conn.Close() + }() + + lastWriteLk := sync.Mutex{} + lastWrite := time.Now() + + // Start a goroutine to ping the client every 30 seconds to check if it's + // still alive. If the client doesn't respond to a ping within 5 seconds, + // we'll close the connection and teardown the consumer. + go func() { + ticker := time.NewTicker(30 * time.Second) + defer ticker.Stop() + + for { + select { + case <-ticker.C: + lastWriteLk.Lock() + lw := lastWrite + lastWriteLk.Unlock() + + if time.Since(lw) < 30*time.Second { + continue + } + + if err := conn.WriteControl(websocket.PingMessage, []byte{}, time.Now().Add(5*time.Second)); err != nil { + r.Logger.Warn("failed to ping client", "err", err) + cancel() + return + } + case <-ctx.Done(): + return + } + } + }() + + conn.SetPingHandler(func(message string) error { + err := conn.WriteControl(websocket.PongMessage, []byte(message), time.Now().Add(time.Second*60)) + if err == websocket.ErrCloseSent { + return nil + } + return err + }) + + // Start a goroutine to read messages from the client and discard them. + go func() { + for { + _, _, err := conn.ReadMessage() + if err != nil { + r.Logger.Warn("failed to read message from client", "err", err) + cancel() + return + } + } + }() + + ident := realIP + "-" + req.UserAgent() + + evts, cleanup, err := r.Events.Subscribe(ctx, ident, func(evt *stream.XRPCStreamEvent) bool { return true }, since) + if err != nil { + return err + } + defer cleanup() + + // Keep track of the consumer for metrics and admin endpoints + consumer := SocketConsumer{ + RemoteAddr: realIP, + UserAgent: req.UserAgent(), + ConnectedAt: time.Now(), + } + sentCounter := eventsSentCounter.WithLabelValues(consumer.RemoteAddr, consumer.UserAgent) + consumer.EventsSent = sentCounter + + consumerID := r.registerConsumer(&consumer) + defer r.cleanupConsumer(consumerID) + + logger := r.Logger.With( + "consumer_id", consumerID, + "remote_addr", consumer.RemoteAddr, + "user_agent", consumer.UserAgent, + ) + + logger.Info("new consumer", "cursor", since) + + for { + select { + case evt, ok := <-evts: + if !ok { + logger.Error("event stream closed unexpectedly") + return nil + } + + wc, err := conn.NextWriter(websocket.BinaryMessage) + if err != nil { + logger.Error("failed to get next writer", "err", err) + return err + } + + if evt.Preserialized != nil { + _, err = wc.Write(evt.Preserialized) + } else { + err = evt.Serialize(wc) + } + if err != nil { + return fmt.Errorf("failed to write event: %w", err) + } + + if err := wc.Close(); err != nil { + logger.Warn("failed to flush-close our event write", "err", err) + return nil + } + + lastWriteLk.Lock() + lastWrite = time.Now() + lastWriteLk.Unlock() + sentCounter.Inc() + case <-ctx.Done(): + return nil + } + } +} + +type ConsumerInfo struct { + ID uint64 `json:"id"` + RemoteAddr string `json:"remote_addr"` + UserAgent string `json:"user_agent"` + EventsConsumed uint64 `json:"events_consumed"` + ConnectedAt time.Time `json:"connected_at"` +} + +func (r *Relay) ListConsumers() []ConsumerInfo { + r.consumersLk.RLock() + defer r.consumersLk.RUnlock() + + info := make([]ConsumerInfo, 0, len(r.consumers)) + for id, c := range r.consumers { + var m = &dto.Metric{} + if err := c.EventsSent.Write(m); err != nil { + continue + } + info = append(info, ConsumerInfo{ + ID: id, + RemoteAddr: c.RemoteAddr, + UserAgent: c.UserAgent, + EventsConsumed: uint64(m.Counter.GetValue()), + ConnectedAt: c.ConnectedAt, + }) + } + return info +} diff --git a/knotmirror/relay/ingest.go b/knotmirror/relay/ingest.go new file mode 100644 index 00000000..85823623 --- /dev/null +++ b/knotmirror/relay/ingest.go @@ -0,0 +1,61 @@ +package relay + +import ( + "context" + "fmt" + "log/slog" + "time" + + "github.com/bluesky-social/indigo/atproto/identity" + "tangled.org/core/knotmirror/stream" +) + +// This callback function gets called by Slurper on every upstream repo stream message from any host. +// +// Messages are processed in-order for a single account on a single host; but may be concurrent or out-of-order for the same account *across* hosts (eg, during account migration or a conflict) +func (r *Relay) processRepoEvent(ctx context.Context, evt stream.EventStreamMessage, hostname string, hostID uint64) error { + ctx, span := tracer.Start(ctx, "processRepoEvent") + defer span.End() + + start := time.Now() + defer func() { + eventsHandleDuration.WithLabelValues(hostname).Observe(time.Since(start).Seconds()) + }() + + EventsReceivedCounter.WithLabelValues(hostname).Add(1) + + // processXXXEvent + logger := r.Logger.With("seq", evt.Sequence(), "host", hostname) + + acc, ident, err := r.preProcessEvent(ctx, evt.DID(), hostname, hostID, logger) + if err != nil { + return err + } + + // verify that the account has active status + if err := r.EnsureAccountActive(ctx, acc); err != nil { + logger.Info("dropping message for inactive account", "status", acc.Status, "upstreamStatus", acc.UpstreamStatus) + eventsWarningsCounter.WithLabelValues(hostname, "inactive-account").Add(1) + return nil + } + + // TODO(git_sync): verify event payload once we start signing events + // 1. fast check for stale revision + // 2. verify event payload + _ = ident + + err = r.Events.AddEvent(ctx, &stream.XRPCStreamEvent{ + Event: evt, + PrivUid: acc.UID, + }) + if err != nil { + logger.Error("failed to broadcast event", "error", err) + return fmt.Errorf("failed to broadcast #commit event: %w", err) + } + + return nil +} + +func (r *Relay) preProcessEvent(ctx context.Context, didStr string, hostname string, hostID uint64, logger *slog.Logger) (*models.Account, *identity.Identity, error) { + panic("unimplemented") +} diff --git a/knotmirror/relay/metrics.go b/knotmirror/relay/metrics.go new file mode 100644 index 00000000..ffb077a6 --- /dev/null +++ b/knotmirror/relay/metrics.go @@ -0,0 +1,49 @@ +package relay + +import ( + "github.com/prometheus/client_golang/prometheus" + "github.com/prometheus/client_golang/prometheus/promauto" +) + +// TODO: expose an accessor instead of exporting +var EventsReceivedCounter = promauto.NewCounterVec(prometheus.CounterOpts{ + Name: "events_received_counter", + Help: "The total number of events received", +}, []string{"pds"}) + +var eventsWarningsCounter = promauto.NewCounterVec(prometheus.CounterOpts{ + Name: "events_warn_counter", + Help: "Events received with warnings", +}, []string{"pds", "warn"}) + +var eventsHandleDuration = promauto.NewHistogramVec(prometheus.HistogramOpts{ + Name: "events_handle_duration", + Help: "A histogram of handleFedEvent latencies", + Buckets: prometheus.ExponentialBuckets(0.001, 2, 15), +}, []string{"pds"}) + +var repoCommitsReceivedCounter = promauto.NewCounterVec(prometheus.CounterOpts{ + Name: "repo_commits_received_counter", + Help: "The total number of commit events received", +}, []string{"pds"}) +var repoSyncReceivedCounter = promauto.NewCounterVec(prometheus.CounterOpts{ + Name: "repo_sync_received_counter", + Help: "The total number of sync events received", +}, []string{"pds"}) + +var eventsSentCounter = promauto.NewCounterVec(prometheus.CounterOpts{ + Name: "events_sent_counter", + Help: "The total number of events sent to consumers", +}, []string{"remote_addr", "user_agent"}) + +/* NOTE: not implemented in this version of relay +var externalUserCreationAttempts = promauto.NewCounter(prometheus.CounterOpts{ + Name: "relay_external_user_creation_attempts", + Help: "The total number of external users created", +}) +*/ + +var newUsersDiscovered = promauto.NewCounter(prometheus.CounterOpts{ + Name: "relay_new_users_discovered", + Help: "The total number of new users discovered directly from the firehose (not from refs)", +}) diff --git a/knotmirror/relay/relay.go b/knotmirror/relay/relay.go new file mode 100644 index 00000000..beeb6381 --- /dev/null +++ b/knotmirror/relay/relay.go @@ -0,0 +1,88 @@ +package relay + +import ( + "log/slog" + "sync" + "time" + + "tangled.org/core/knotmirror/stream/eventmgr" + + "github.com/RussellLuo/slidingwindow" + "github.com/bluesky-social/indigo/atproto/identity" + lru "github.com/hashicorp/golang-lru/v2" + "go.opentelemetry.io/otel" +) + +var tracer = otel.Tracer("relay") + +type Relay struct { + db *gorm.DB + Dir identity.Directory + Logger *slog.Logger + Slurper *Slurper + Events *eventmgr.EventManager + HostChecker HostChecker + Config RelayConfig + + // Management of Socket Consumers + consumersLk sync.RWMutex + nextConsumerID uint64 + consumers map[uint64]*SocketConsumer + + // Account cache + accountCache *lru.Cache[string, *models.Account] + + HostPerDayLimiter *slidingwindow.Limiter +} + +func NewRelay(db *gorm.DB, evtman *eventmgr.EventManager, dir identity.Directory, config *RelayConfig) (*Relay, error) { + + if config == nil { + config = DefaultRelayConfig() + } + + uc, _ := lru.New[string, *models.Account](2_000_000) + + hc := NewHostClient(config.UserAgent) + + r := &Relay{ + db: db, + Dir: dir, + Logger: slog.Default().With("system", "relay"), + Events: evtman, + HostChecker: hc, + Config: *config, + + consumers: make(map[uint64]*SocketConsumer), + + accountCache: uc, + + HostPerDayLimiter: perDayLimiter(config.HostPerDayLimit), + } + + if err := r.MigrateDatabase(); err != nil { + return nil, err + } + + slurpConfig := DefaultSlurperConfig() + slurpConfig.ConcurrencyPerHost = config.ConcurrencyPerHost + + // register callbacks to persist cursors and host state in database + slurpConfig.PersistCursorCallback = r.PersistHostCursors + slurpConfig.PersistHostStatusCallback = r.UpdateHostStatus + + s, err := NewSlurper(r.processRepoEvent, slurpConfig) + if err != nil { + return nil, err + } + r.Slurper = s + + return r, nil +} + +func perDayLimiter(count int64) *slidingwindow.Limiter { + lim, _ := slidingwindow.NewLimiter(time.Hour*24, count, func() (slidingwindow.Window, slidingwindow.StopFunc) { + return slidingwindow.NewLocalWindow() + }) + return lim +} -- 2.51.2