Version 1.0.0
This commit is contained in:
@@ -2,6 +2,9 @@ package notify
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"fmt"
|
||||
"math/rand"
|
||||
"net"
|
||||
"strings"
|
||||
"time"
|
||||
@@ -26,13 +29,16 @@ type StarNotifyC struct {
|
||||
// Queue 是用来处理收发信息的简单消息队列
|
||||
Queue *starainrt.StarQueue
|
||||
// Online 当前链接是否处于活跃状态
|
||||
Online bool
|
||||
Online bool
|
||||
lockPool map[string]CMsg
|
||||
}
|
||||
|
||||
// CMsg 指明当前客户端被通知的关键字
|
||||
type CMsg struct {
|
||||
Key string
|
||||
Value string
|
||||
mode string
|
||||
wait chan int
|
||||
}
|
||||
|
||||
func (star *StarNotifyC) starinitc() {
|
||||
@@ -43,6 +49,7 @@ func (star *StarNotifyC) starinitc() {
|
||||
star.Stop = make(chan int, 5)
|
||||
star.clientSign = make(map[string]chan string)
|
||||
star.Online = false
|
||||
star.lockPool = make(map[string]CMsg)
|
||||
star.Queue.RestoreDuration(time.Second * 2)
|
||||
}
|
||||
|
||||
@@ -107,13 +114,84 @@ func NewNotifyC(netype, value string) (*StarNotifyC, error) {
|
||||
|
||||
// Send 用于向Server端发送数据
|
||||
func (star *StarNotifyC) Send(name string) error {
|
||||
_, err := star.Connc.Write(star.Queue.BuildMessage([]byte(name)))
|
||||
return err
|
||||
return star.SendValue(name, "")
|
||||
}
|
||||
|
||||
// SendValue 用于向Server端发送key-value类型数据
|
||||
func (star *StarNotifyC) SendValue(name, value string) error {
|
||||
_, err := star.Connc.Write(star.Queue.BuildMessage([]byte(name + "||" + value)))
|
||||
var err error
|
||||
var key []byte
|
||||
for _, v := range []byte(name) {
|
||||
if v == byte(124) || v == byte(92) {
|
||||
key = append(key, byte(92))
|
||||
}
|
||||
key = append(key, v)
|
||||
}
|
||||
_, err = star.Connc.Write(star.Queue.BuildMessage([]byte("pa" + "||" + string(key) + "||" + value)))
|
||||
return err
|
||||
}
|
||||
|
||||
func (star *StarNotifyC) trim(name string) string {
|
||||
var slash bool = false
|
||||
var key []byte
|
||||
for _, v := range []byte(name) {
|
||||
if v == byte(92) && !slash {
|
||||
slash = true
|
||||
continue
|
||||
}
|
||||
slash = false
|
||||
key = append(key, v)
|
||||
}
|
||||
return string(key)
|
||||
}
|
||||
|
||||
// SendValueWait 用于向Server端发送key-value类型数据并等待结果返回,此结果不会通过标准返回流程处理
|
||||
func (star *StarNotifyC) SendValueWait(name, value string, tmout time.Duration) (CMsg, error) {
|
||||
var err error
|
||||
var tmceed <-chan time.Time
|
||||
if star.UseChannel {
|
||||
return CMsg{}, errors.New("Do Not Use UseChannel Mode!")
|
||||
}
|
||||
rand.Seed(time.Now().UnixNano())
|
||||
mode := "cr" + fmt.Sprintf("%05d", rand.Intn(99999))
|
||||
var key []byte
|
||||
for _, v := range []byte(name) {
|
||||
if v == byte(124) || v == byte(92) {
|
||||
key = append(key, byte(92))
|
||||
}
|
||||
key = append(key, v)
|
||||
}
|
||||
_, err = star.Connc.Write(star.Queue.BuildMessage([]byte(mode + "||" + string(key) + "||" + value)))
|
||||
if err != nil {
|
||||
return CMsg{}, err
|
||||
}
|
||||
if int64(tmout) > 0 {
|
||||
tmceed = time.After(tmout)
|
||||
}
|
||||
var source CMsg
|
||||
source.wait = make(chan int, 2)
|
||||
star.lockPool[mode] = source
|
||||
select {
|
||||
case <-source.wait:
|
||||
res := star.lockPool[mode]
|
||||
delete(star.lockPool, mode)
|
||||
return res, nil
|
||||
case <-tmceed:
|
||||
return CMsg{}, errors.New("Time Exceed")
|
||||
}
|
||||
}
|
||||
|
||||
// ReplyMsg 用于向Server端Reply信息
|
||||
func (star *StarNotifyC) ReplyMsg(data CMsg, name, value string) error {
|
||||
var err error
|
||||
var key []byte
|
||||
for _, v := range []byte(name) {
|
||||
if v == byte(124) || v == byte(92) {
|
||||
key = append(key, byte(92))
|
||||
}
|
||||
key = append(key, v)
|
||||
}
|
||||
_, err = star.Connc.Write(star.Queue.BuildMessage([]byte(data.mode + "||" + string(key) + "||" + value)))
|
||||
return err
|
||||
}
|
||||
|
||||
@@ -134,19 +212,38 @@ func (star *StarNotifyC) cnotify() {
|
||||
star.Online = false
|
||||
return
|
||||
}
|
||||
strs := strings.SplitN(string(data.Msg), "||", 2)
|
||||
if len(strs) < 2 {
|
||||
strs := strings.SplitN(string(data.Msg), "||", 3)
|
||||
if len(strs) < 3 {
|
||||
continue
|
||||
}
|
||||
strs[1] = star.trim(strs[1])
|
||||
if star.UseChannel {
|
||||
go star.store(strs[0], strs[1])
|
||||
go star.store(strs[1], strs[2])
|
||||
} else {
|
||||
key, value := strs[0], strs[1]
|
||||
if msg, ok := star.FuncLists[key]; ok {
|
||||
go msg(CMsg{key, value})
|
||||
mode, key, value := strs[0], strs[1], strs[2]
|
||||
if mode[0:2] != "cr" {
|
||||
if msg, ok := star.FuncLists[key]; ok {
|
||||
go msg(CMsg{key, value, mode, nil})
|
||||
} else {
|
||||
if star.defaultFunc != nil {
|
||||
go star.defaultFunc(CMsg{key, value, mode, nil})
|
||||
}
|
||||
}
|
||||
} else {
|
||||
if star.defaultFunc != nil {
|
||||
go star.defaultFunc(CMsg{key, value})
|
||||
if sa, ok := star.lockPool[mode]; ok {
|
||||
sa.Key = key
|
||||
sa.Value = value
|
||||
sa.mode = mode
|
||||
star.lockPool[mode] = sa
|
||||
sa.wait <- 1
|
||||
} else {
|
||||
if msg, ok := star.FuncLists[key]; ok {
|
||||
go msg(CMsg{key, value, mode, nil})
|
||||
} else {
|
||||
if star.defaultFunc != nil {
|
||||
go star.defaultFunc(CMsg{key, value, mode, nil})
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -160,6 +257,8 @@ func (star *StarNotifyC) ClientStop() {
|
||||
}
|
||||
star.cancel()
|
||||
star.Stop <- 1
|
||||
star.Stop <- 1
|
||||
star.Stop <- 1
|
||||
}
|
||||
|
||||
// SetNotify 用于设置关键词的调用函数
|
||||
@@ -168,6 +267,6 @@ func (star *StarNotifyC) SetNotify(name string, data func(CMsg)) {
|
||||
}
|
||||
|
||||
// SetDefaultNotify 用于设置默认关键词的调用函数
|
||||
func (star *StarNotifyC) SetDefaultNotify(name string, data func(CMsg)) {
|
||||
func (star *StarNotifyC) SetDefaultNotify(data func(CMsg)) {
|
||||
star.defaultFunc = data
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user