ing
This commit is contained in:
parent
435f308cda
commit
d5249ae7e8
|
@ -2,7 +2,6 @@ package redis
|
||||||
|
|
||||||
import (
|
import (
|
||||||
"context"
|
"context"
|
||||||
"log"
|
|
||||||
|
|
||||||
channelUtil "git.loafle.net/commons_go/util/channel"
|
channelUtil "git.loafle.net/commons_go/util/channel"
|
||||||
ofSubscriber "git.loafle.net/overflow/overflow_subscriber"
|
ofSubscriber "git.loafle.net/overflow/overflow_subscriber"
|
||||||
|
@ -25,6 +24,7 @@ type subscriber struct {
|
||||||
subListeners map[string]ofSubscriber.OnSubscribeFunc
|
subListeners map[string]ofSubscriber.OnSubscribeFunc
|
||||||
isListenSubscriptions bool
|
isListenSubscriptions bool
|
||||||
subCh chan subscribeChannelAction
|
subCh chan subscribeChannelAction
|
||||||
|
isRunning bool
|
||||||
}
|
}
|
||||||
|
|
||||||
func New(ctx context.Context, conn redis.Conn) Subscriber {
|
func New(ctx context.Context, conn redis.Conn) Subscriber {
|
||||||
|
@ -38,6 +38,8 @@ func New(ctx context.Context, conn redis.Conn) Subscriber {
|
||||||
|
|
||||||
go n.listen()
|
go n.listen()
|
||||||
|
|
||||||
|
n.isRunning = true
|
||||||
|
|
||||||
return n
|
return n
|
||||||
}
|
}
|
||||||
|
|
||||||
|
@ -48,9 +50,7 @@ func (n *subscriber) listen() {
|
||||||
switch sa.Type {
|
switch sa.Type {
|
||||||
case channelUtil.ActionTypeCreate:
|
case channelUtil.ActionTypeCreate:
|
||||||
_, ok := n.subListeners[sa.channel]
|
_, ok := n.subListeners[sa.channel]
|
||||||
if ok {
|
if !ok {
|
||||||
log.Fatalf("Subscriber: Subscription of channel[%s] is already exist", sa.channel)
|
|
||||||
} else {
|
|
||||||
n.subListeners[sa.channel] = sa.cb
|
n.subListeners[sa.channel] = sa.cb
|
||||||
n.conn.Subscribe(sa.channel)
|
n.conn.Subscribe(sa.channel)
|
||||||
n.listenSubscriptions()
|
n.listenSubscriptions()
|
||||||
|
@ -61,19 +61,21 @@ func (n *subscriber) listen() {
|
||||||
if ok {
|
if ok {
|
||||||
n.conn.Unsubscribe(sa.channel)
|
n.conn.Unsubscribe(sa.channel)
|
||||||
delete(n.subListeners, sa.channel)
|
delete(n.subListeners, sa.channel)
|
||||||
} else {
|
|
||||||
log.Fatalf("Subscriber: Subscription of channel[%s] is not exist", sa.channel)
|
|
||||||
}
|
}
|
||||||
break
|
break
|
||||||
}
|
}
|
||||||
case <-n.ctx.Done():
|
case <-n.ctx.Done():
|
||||||
log.Println("redis subscriber: Context Done")
|
n.destroy()
|
||||||
n.conn.Close()
|
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func (n *subscriber) destroy() {
|
||||||
|
n.isRunning = false
|
||||||
|
n.conn.Close()
|
||||||
|
}
|
||||||
|
|
||||||
func (n *subscriber) listenSubscriptions() {
|
func (n *subscriber) listenSubscriptions() {
|
||||||
if n.isListenSubscriptions {
|
if n.isListenSubscriptions {
|
||||||
return
|
return
|
||||||
|
@ -86,10 +88,11 @@ func (n *subscriber) listenSubscriptions() {
|
||||||
if cb, ok := n.subListeners[v.Channel]; ok {
|
if cb, ok := n.subListeners[v.Channel]; ok {
|
||||||
cb(v.Channel, string(v.Data))
|
cb(v.Channel, string(v.Data))
|
||||||
}
|
}
|
||||||
|
break
|
||||||
case redis.Subscription:
|
case redis.Subscription:
|
||||||
log.Printf("subscription message: %s: %s %d\n", v.Channel, v.Kind, v.Count)
|
break
|
||||||
case error:
|
case error:
|
||||||
log.Println("error pub/sub, delivery has stopped")
|
n.destroy()
|
||||||
return
|
return
|
||||||
default:
|
default:
|
||||||
}
|
}
|
||||||
|
|
Loading…
Reference in New Issue
Block a user