2021-11-01 10:42:55 -04:00
|
|
|
package relay
|
2021-04-14 22:17:12 -04:00
|
|
|
|
|
|
|
import (
|
|
|
|
"sync"
|
|
|
|
|
2022-11-09 15:53:01 -04:00
|
|
|
"github.com/waku-org/go-waku/waku/v2/protocol"
|
2021-04-14 22:17:12 -04:00
|
|
|
)
|
|
|
|
|
2021-10-09 14:18:53 -04:00
|
|
|
// Subscription handles the subscrition to a particular pubsub topic
|
2021-04-14 22:17:12 -04:00
|
|
|
type Subscription struct {
|
2022-02-23 11:08:27 -04:00
|
|
|
sync.RWMutex
|
|
|
|
|
2021-10-09 14:18:53 -04:00
|
|
|
// C is channel used for receiving envelopes
|
2021-04-22 14:49:52 -04:00
|
|
|
C chan *protocol.Envelope
|
|
|
|
|
2021-04-14 22:17:12 -04:00
|
|
|
closed bool
|
2021-11-01 10:42:55 -04:00
|
|
|
once sync.Once
|
2021-04-14 22:17:12 -04:00
|
|
|
quit chan struct{}
|
|
|
|
}
|
|
|
|
|
2021-10-09 14:18:53 -04:00
|
|
|
// Unsubscribe will close a subscription from a pubsub topic. Will close the message channel
|
2021-04-14 22:17:12 -04:00
|
|
|
func (subs *Subscription) Unsubscribe() {
|
2021-11-01 10:42:55 -04:00
|
|
|
subs.once.Do(func() {
|
|
|
|
close(subs.quit)
|
|
|
|
})
|
2022-02-23 11:08:27 -04:00
|
|
|
|
2021-04-14 22:17:12 -04:00
|
|
|
}
|
|
|
|
|
2021-10-09 14:18:53 -04:00
|
|
|
// IsClosed determine whether a Subscription is still open for receiving messages
|
2021-04-14 22:17:12 -04:00
|
|
|
func (subs *Subscription) IsClosed() bool {
|
2022-02-23 11:08:27 -04:00
|
|
|
subs.RLock()
|
|
|
|
defer subs.RUnlock()
|
2021-04-14 22:17:12 -04:00
|
|
|
return subs.closed
|
|
|
|
}
|