1
0
Fork 0
mirror of https://github.com/eosswedenorg/thalos synced 2026-06-16 04:24:56 +02:00

api/redis/subscriber.go: adding some comments.

This commit is contained in:
Henrik Hautakoski 2024-02-04 14:48:54 +01:00
parent 35a9706954
commit 728b03422f

View file

@ -11,8 +11,10 @@ import (
)
type Subscriber struct {
sub *redis.PubSub
ctx context.Context
sub *redis.PubSub
ctx context.Context
// Mutex for channels map.
mu sync.RWMutex
timeout time.Duration
channels map[string]chan []byte
@ -94,11 +96,15 @@ func (s *Subscriber) Read(channel api.Channel) ([]byte, error) {
}
func (s *Subscriber) Close() error {
// Close redis pubsub.
err := s.sub.Close()
// Close all go channels, this will make Read() unblock.
for _, ch := range s.channels {
close(ch)
}
// Clear the channel map of old channels.
s.mu.Lock()
s.channels = make(map[string]chan []byte)
s.mu.Unlock()