Files
Gitea/services/pubsub/redis.go
T
silverwind 7efcb8d6ca test(pubsub): stop racing the Redis SUBSCRIBE ack (#38661)
`RedisBroker.Subscribe` returns before the server acks `SUBSCRIBE`, so a
publish right after it can be dropped, making
`TestRedisBroker/CrossBroker` fail intermittently on loaded CI runners
([example](https://github.com/go-gitea/gitea/actions/runs/30257188479/job/89948425118)).

Each scenario now uses its own topic and waits for `PUBSUB NUMSUB`
before publishing. `MemoryBroker` registers synchronously and skips the
wait.
2026-07-27 16:12:06 +02:00

193 lines
5.4 KiB
Go

// 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
}