// Copyright 2026 The Gitea Authors. All rights reserved. // SPDX-License-Identifier: MIT package pubsub import ( "context" "sync" "time" "gitea.dev/modules/graceful" "gitea.dev/modules/log" "gitea.dev/modules/nosql" "gitea.dev/modules/util" "github.com/redis/go-redis/v9" ) const ( redisPingTimeout = 3 * time.Second redisPingRetries = 10 redisPingRetryDelay = time.Second redisPublishTimeout = 2 * time.Second ) // RedisBroker fans out across processes via Redis pub/sub. Each topic is // backed by a single Redis SUBSCRIBE shared between local subscribers; the // last local Unsubscribe tears the Redis subscription down. type RedisBroker struct { client redis.UniversalClient mu sync.RWMutex topics map[string]*redisTopic } type redisTopic struct { ps *redis.PubSub subs []*redisSub cancel context.CancelFunc } // redisSub pairs a delivery channel with the once that guards its close, so // either cancel() or readLoop's error-exit path can safely close it. type redisSub struct { ch chan []byte once sync.Once } func (s *redisSub) close() { s.once.Do(func() { close(s.ch) }) } var _ Broker = (*RedisBroker)(nil) func redisChannelForTopic(s string) string { return "gitea-ws-topic:" + s } func NewRedisBroker(connStr string) (*RedisBroker, error) { client := nosql.GetManager().GetRedisClient(connStr) // context.Background not graceful.ShutdownContext: shutdown ctx may not be initialized at boot. // Retry to ride out docker-compose start-order races (matches modules/queue). var err error for range redisPingRetries { pingCtx, cancel := context.WithTimeout(context.Background(), redisPingTimeout) err = client.Ping(pingCtx).Err() cancel() if err == nil { break } log.Warn("pubsub redis: not ready, retrying in 1s: %v", err) time.Sleep(redisPingRetryDelay) } if err != nil { return nil, err } return &RedisBroker{ client: client, topics: make(map[string]*redisTopic), }, nil } func (b *RedisBroker) Subscribe(topic string) (<-chan []byte, func()) { sub := &redisSub{ch: make(chan []byte, subChanBuffer)} // Fast path: topic already has a Redis subscription, just attach locally. b.mu.Lock() if state, exists := b.topics[topic]; exists { state.subs = append(state.subs, sub) b.mu.Unlock() return sub.ch, b.makeCancel(topic, sub) } b.mu.Unlock() // Slow path: create the Redis subscription outside the broker mutex so // other Subscribe/cancel calls aren't blocked on the network round-trip. // graceful.ShutdownContext so the reader loop dies cleanly on Gitea // shutdown even if every local subscriber has already cancelled. // readLoop consumes the SUBSCRIBE ack; don't wait for it here, that would put // a Redis round-trip in the WebSocket handshake to close a sub-millisecond // window in which a publish is missed - harmless, since a client receives // nothing at all until its handshake completes. ctx, cancelCtx := context.WithCancel(graceful.GetManager().ShutdownContext()) ps := b.client.Subscribe(ctx, redisChannelForTopic(topic)) b.mu.Lock() if existing, exists := b.topics[topic]; exists { // Another goroutine won the create race; merge into theirs and discard ours. existing.subs = append(existing.subs, sub) b.mu.Unlock() cancelCtx() _ = ps.Close() return sub.ch, b.makeCancel(topic, sub) } b.topics[topic] = &redisTopic{ps: ps, cancel: cancelCtx, subs: []*redisSub{sub}} b.mu.Unlock() go b.readLoop(ctx, topic, ps) return sub.ch, b.makeCancel(topic, sub) } func (b *RedisBroker) makeCancel(topic string, sub *redisSub) func() { return func() { b.mu.Lock() state, ok := b.topics[topic] if !ok { b.mu.Unlock() sub.close() return } state.subs = util.SliceRemoveAll(state.subs, sub) if len(state.subs) == 0 { state.cancel() _ = state.ps.Close() delete(b.topics, topic) } b.mu.Unlock() sub.close() } } func (b *RedisBroker) readLoop(ctx context.Context, topic string, ps *redis.PubSub) { for { msg, err := ps.ReceiveMessage(ctx) if err != nil { if ctx.Err() != nil { return } // Transport blip: tear the topic down so a fresh Subscribe rebuilds it. // Closing each subscriber's channel wakes the WebSocket handler, which // will reconnect and re-Subscribe. b.mu.Lock() if cur, ok := b.topics[topic]; ok && cur.ps == ps { for _, s := range cur.subs { s.close() } cur.cancel() _ = cur.ps.Close() delete(b.topics, topic) } b.mu.Unlock() log.Trace("pubsub redis: receive on %q: %v", topic, err) return } payload := []byte(msg.Payload) b.mu.RLock() state, ok := b.topics[topic] if !ok { b.mu.RUnlock() return } for _, s := range state.subs { select { case s.ch <- payload: default: log.Trace("pubsub redis: dropping message on topic %q — subscriber channel full", topic) } } b.mu.RUnlock() } } func (b *RedisBroker) Publish(topic string, msg []byte) { ctx, cancel := context.WithTimeout(graceful.GetManager().HammerContext(), redisPublishTimeout) defer cancel() if err := b.client.Publish(ctx, redisChannelForTopic(topic), msg).Err(); err != nil { log.Error("pubsub redis: publish to %q: %v", topic, err) } } // HasTopicSubscribers conservatively returns true: cross-process subscriber // discovery via PUBSUB NUMSUB is per-node and would silently miss subscribers // in cluster mode. Publishers do the upstream lookup unconditionally. func (b *RedisBroker) HasTopicSubscribers(topic string) bool { return true }