2018-09-07 12:32:26 +00:00
|
|
|
package subscription
|
|
|
|
|
|
|
|
import (
|
2018-09-24 14:13:57 +00:00
|
|
|
"errors"
|
2023-09-19 02:01:38 +00:00
|
|
|
"gitea.watsonlabs.net/watsonb8/fayec/message"
|
2018-09-12 07:42:45 +00:00
|
|
|
"regexp"
|
2018-09-07 12:32:26 +00:00
|
|
|
)
|
|
|
|
|
2018-09-24 14:13:57 +00:00
|
|
|
var ErrInvalidChannelName = errors.New("invalid channel channel")
|
2018-09-07 13:00:35 +00:00
|
|
|
|
2018-09-24 14:13:57 +00:00
|
|
|
type Unsubscriber func(subscription *Subscription) error
|
2018-09-07 13:00:35 +00:00
|
|
|
|
2018-09-07 12:32:26 +00:00
|
|
|
type Subscription struct {
|
2023-09-20 01:40:52 +00:00
|
|
|
channel string
|
|
|
|
authToken string
|
|
|
|
unsub Unsubscriber
|
|
|
|
msgCh chan *message.Message
|
2018-09-07 12:32:26 +00:00
|
|
|
}
|
|
|
|
|
2023-09-19 01:58:15 +00:00
|
|
|
// todo error
|
2023-09-20 01:40:52 +00:00
|
|
|
func NewSubscription(chanel string, unsub Unsubscriber, authToken string, msgCh chan *message.Message) (*Subscription, error) {
|
2018-09-24 14:13:57 +00:00
|
|
|
if !IsValidSubscriptionName(chanel) {
|
|
|
|
return nil, ErrInvalidChannelName
|
|
|
|
}
|
2018-09-07 12:32:26 +00:00
|
|
|
return &Subscription{
|
2023-09-20 01:40:52 +00:00
|
|
|
channel: chanel,
|
|
|
|
authToken: authToken,
|
|
|
|
unsub: unsub,
|
|
|
|
msgCh: msgCh,
|
2018-09-24 14:13:57 +00:00
|
|
|
}, nil
|
2018-09-07 12:32:26 +00:00
|
|
|
}
|
|
|
|
|
2018-09-24 14:13:57 +00:00
|
|
|
func (s *Subscription) OnMessage(onMessage func(channel string, msg message.Data)) error {
|
2018-09-07 12:32:26 +00:00
|
|
|
var inMsg *message.Message
|
2023-09-20 01:40:52 +00:00
|
|
|
for inMsg = range s.msgCh {
|
|
|
|
if inMsg.GetError() != nil {
|
|
|
|
return inMsg.GetError()
|
2018-09-07 12:32:26 +00:00
|
|
|
}
|
2023-09-20 01:40:52 +00:00
|
|
|
onMessage(inMsg.Channel, inMsg.Data)
|
|
|
|
}
|
2018-09-07 12:32:26 +00:00
|
|
|
return nil
|
|
|
|
}
|
|
|
|
|
|
|
|
func (s *Subscription) MsgChannel() chan *message.Message {
|
|
|
|
return s.msgCh
|
|
|
|
}
|
|
|
|
|
2018-09-24 14:13:57 +00:00
|
|
|
func (s *Subscription) Name() string {
|
2018-09-07 12:32:26 +00:00
|
|
|
return s.channel
|
|
|
|
}
|
|
|
|
|
2023-09-19 01:58:15 +00:00
|
|
|
// Unsubscribe ...
|
2018-09-07 12:32:26 +00:00
|
|
|
func (s *Subscription) Unsubscribe() error {
|
2018-09-07 13:00:35 +00:00
|
|
|
return s.unsub(s)
|
|
|
|
}
|
|
|
|
|
2023-09-20 01:40:52 +00:00
|
|
|
func (s *Subscription) AuthToken() string {
|
|
|
|
return s.authToken
|
|
|
|
}
|
|
|
|
|
2023-09-19 01:58:15 +00:00
|
|
|
// validChannelName channel specifies is the channel is in the format /foo/432/bar
|
2018-09-12 07:42:45 +00:00
|
|
|
var validChannelName = regexp.MustCompile(`^\/(((([a-z]|[A-Z])|[0-9])|(\-|\_|\!|\~|\(|\)|\$|\@)))+(\/(((([a-z]|[A-Z])|[0-9])|(\-|\_|\!|\~|\(|\)|\$|\@)))+)*$`)
|
2018-09-24 14:13:57 +00:00
|
|
|
|
2018-09-12 07:42:45 +00:00
|
|
|
var validChannelPattern = regexp.MustCompile(`^(\/(((([a-z]|[A-Z])|[0-9])|(\-|\_|\!|\~|\(|\)|\$|\@)))+)*\/\*{1,2}$`)
|
|
|
|
|
2018-09-24 14:13:57 +00:00
|
|
|
func IsValidSubscriptionName(channel string) bool {
|
2018-09-12 07:42:45 +00:00
|
|
|
return validChannelName.MatchString(channel) || validChannelPattern.MatchString(channel)
|
|
|
|
}
|
2018-09-24 14:13:57 +00:00
|
|
|
|
2023-09-19 01:58:15 +00:00
|
|
|
// isValidPublishName
|
2018-09-24 14:13:57 +00:00
|
|
|
func IsValidPublishName(channel string) bool {
|
|
|
|
return validChannelName.MatchString(channel)
|
|
|
|
}
|