ing
This commit is contained in:
10
notification/notify.go
Normal file
10
notification/notify.go
Normal file
@@ -0,0 +1,10 @@
|
||||
package notification
|
||||
|
||||
type (
|
||||
OnSubscribeFunc func(channel string, payload string)
|
||||
)
|
||||
|
||||
type Notifier interface {
|
||||
Subscribe(channel string, cb OnSubscribeFunc)
|
||||
Unsubscribe(channel string, cb OnSubscribeFunc)
|
||||
}
|
||||
120
notification/redis/notify.go
Normal file
120
notification/redis/notify.go
Normal file
@@ -0,0 +1,120 @@
|
||||
package redis
|
||||
|
||||
import (
|
||||
"context"
|
||||
"log"
|
||||
|
||||
channelUtil "git.loafle.net/commons_go/util/channel"
|
||||
"git.loafle.net/overflow/overflow_gateway_web/notification"
|
||||
"github.com/garyburd/redigo/redis"
|
||||
)
|
||||
|
||||
type subscribeChannelAction struct {
|
||||
channelUtil.Action
|
||||
channel string
|
||||
cb notification.OnSubscribeFunc
|
||||
}
|
||||
|
||||
type Notifier interface {
|
||||
notification.Notifier
|
||||
}
|
||||
|
||||
type notifier struct {
|
||||
ctx context.Context
|
||||
conn redis.PubSubConn
|
||||
subListeners map[string]notification.OnSubscribeFunc
|
||||
isListenSubscriptions bool
|
||||
subCh chan subscribeChannelAction
|
||||
}
|
||||
|
||||
func New(ctx context.Context, conn redis.Conn) Notifier {
|
||||
n := ¬ifier{
|
||||
ctx: ctx,
|
||||
subListeners: make(map[string]notification.OnSubscribeFunc),
|
||||
isListenSubscriptions: false,
|
||||
subCh: make(chan subscribeChannelAction),
|
||||
}
|
||||
n.conn = redis.PubSubConn{Conn: conn}
|
||||
|
||||
go n.listen()
|
||||
|
||||
return n
|
||||
}
|
||||
|
||||
func (n *notifier) listen() {
|
||||
for {
|
||||
select {
|
||||
case sa := <-n.subCh:
|
||||
switch sa.Type {
|
||||
case channelUtil.ActionTypeCreate:
|
||||
_, ok := n.subListeners[sa.channel]
|
||||
if ok {
|
||||
log.Fatalf("notification: Subscription of channel[%s] is already exist", sa.channel)
|
||||
} else {
|
||||
n.subListeners[sa.channel] = sa.cb
|
||||
n.conn.Subscribe(sa.channel)
|
||||
n.listenSubscriptions()
|
||||
}
|
||||
break
|
||||
case channelUtil.ActionTypeDelete:
|
||||
_, ok := n.subListeners[sa.channel]
|
||||
if ok {
|
||||
n.conn.Unsubscribe(sa.channel)
|
||||
delete(n.subListeners, sa.channel)
|
||||
} else {
|
||||
log.Fatalf("notification: Subscription of channel[%s] is not exist", sa.channel)
|
||||
}
|
||||
break
|
||||
}
|
||||
case <-n.ctx.Done():
|
||||
log.Println("redis noti: Context Done")
|
||||
n.conn.Close()
|
||||
return
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func (n *notifier) listenSubscriptions() {
|
||||
if n.isListenSubscriptions {
|
||||
return
|
||||
}
|
||||
|
||||
go func() {
|
||||
for {
|
||||
switch v := n.conn.Receive().(type) {
|
||||
case redis.Message:
|
||||
if cb, ok := n.subListeners[v.Channel]; ok {
|
||||
cb(v.Channel, string(v.Data))
|
||||
}
|
||||
case redis.Subscription:
|
||||
log.Printf("subscription message: %s: %s %d\n", v.Channel, v.Kind, v.Count)
|
||||
case error:
|
||||
log.Println("error pub/sub, delivery has stopped")
|
||||
return
|
||||
default:
|
||||
}
|
||||
}
|
||||
}()
|
||||
|
||||
n.isListenSubscriptions = true
|
||||
}
|
||||
|
||||
func (n *notifier) Subscribe(channel string, cb notification.OnSubscribeFunc) {
|
||||
ca := subscribeChannelAction{
|
||||
channel: channel,
|
||||
cb: cb,
|
||||
}
|
||||
ca.Type = channelUtil.ActionTypeCreate
|
||||
|
||||
n.subCh <- ca
|
||||
}
|
||||
|
||||
func (n *notifier) Unsubscribe(channel string, cb notification.OnSubscribeFunc) {
|
||||
ca := subscribeChannelAction{
|
||||
channel: channel,
|
||||
cb: cb,
|
||||
}
|
||||
ca.Type = channelUtil.ActionTypeDelete
|
||||
|
||||
n.subCh <- ca
|
||||
}
|
||||
Reference in New Issue
Block a user