mirror of
https://github.com/logos-messaging/logos-messaging-go.git
synced 2026-01-08 00:43:10 +00:00
fix: start watchTopicShards once
This commit is contained in:
parent
84a4b1be7a
commit
4be813d71e
@ -348,6 +348,10 @@ func (w *WakuNode) SetRelayShards(rs protocol.RelayShards) error {
|
||||
}
|
||||
|
||||
func (w *WakuNode) watchTopicShards(ctx context.Context) error {
|
||||
if !w.watchingRelayShards.CompareAndSwap(false, true) {
|
||||
return nil
|
||||
}
|
||||
|
||||
evtRelaySubscribed, err := w.Relay().Events().Subscribe(new(relay.EvtRelaySubscribed))
|
||||
if err != nil {
|
||||
return err
|
||||
@ -358,10 +362,13 @@ func (w *WakuNode) watchTopicShards(ctx context.Context) error {
|
||||
return err
|
||||
}
|
||||
|
||||
w.wg.Add(1)
|
||||
|
||||
go func() {
|
||||
defer utils.LogOnPanic()
|
||||
defer evtRelaySubscribed.Close()
|
||||
defer evtRelayUnsubscribed.Close()
|
||||
defer w.wg.Done()
|
||||
|
||||
for {
|
||||
select {
|
||||
|
||||
@ -5,6 +5,7 @@ import (
|
||||
"math/rand"
|
||||
"net"
|
||||
"sync"
|
||||
"sync/atomic"
|
||||
"time"
|
||||
|
||||
backoffv4 "github.com/cenkalti/backoff/v4"
|
||||
@ -122,6 +123,8 @@ type WakuNode struct {
|
||||
storeFactory storeFactory
|
||||
|
||||
peermanager *peermanager.PeerManager
|
||||
|
||||
watchingRelayShards atomic.Bool
|
||||
}
|
||||
|
||||
func defaultStoreFactory(w *WakuNode) legacy_store.Store {
|
||||
|
||||
Loading…
x
Reference in New Issue
Block a user