forked from eolinker/apinto
-
Notifications
You must be signed in to change notification settings - Fork 0
Commit
This commit does not belong to any branch on this repository, and may belong to a fork outside of the repository.
Merge remote-tracking branch 'gitlab/feature/fix_discovery'
- Loading branch information
Showing
78 changed files
with
1,551 additions
and
1,159 deletions.
There are no files selected for viewing
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -1,106 +1,66 @@ | ||
package discovery | ||
|
||
import ( | ||
"github.com/google/uuid" | ||
"github.com/eolinker/eosc/eocontext" | ||
"sync" | ||
"sync/atomic" | ||
) | ||
|
||
type app struct { | ||
id string | ||
nodes map[string]INode | ||
healthChecker IHealthChecker | ||
attrs Attrs | ||
locker sync.RWMutex | ||
container IAppContainer | ||
} | ||
|
||
// Reset 重置app的节点列表 | ||
func (s *app) Reset(nodes Nodes) { | ||
|
||
tmp := make(map[string]INode) | ||
var ( | ||
_ IAppAgent = (*_AppAgent)(nil) | ||
_ IApp = (*_App)(nil) | ||
) | ||
|
||
for _, node := range nodes { | ||
type IAppAgent interface { | ||
reset(nodes []eocontext.INode) | ||
Agent() IApp | ||
} | ||
|
||
if n, has := s.nodes[node.ID()]; has { | ||
n.Leave() | ||
} | ||
tmp[node.ID()] = node | ||
type IApp interface { | ||
Nodes() []eocontext.INode | ||
Close() | ||
} | ||
|
||
} | ||
s.locker.Lock() | ||
s.nodes = tmp | ||
s.locker.Unlock() | ||
type _AppAgent struct { | ||
locker sync.RWMutex | ||
nodes []eocontext.INode | ||
use int64 | ||
} | ||
|
||
// GetAttrs 获取app的属性集合 | ||
func (s *app) GetAttrs() Attrs { | ||
s.locker.RLock() | ||
defer s.locker.RUnlock() | ||
return s.attrs | ||
func (a *_AppAgent) Agent() IApp { | ||
atomic.AddInt64(&a.use, 1) | ||
return &_App{_AppAgent: a, isClose: 0} | ||
} | ||
|
||
// GetAttrByName 通过属性名获取app对应属性 | ||
func (s *app) GetAttrByName(name string) (string, bool) { | ||
s.locker.RLock() | ||
defer s.locker.RUnlock() | ||
attr, ok := s.attrs[name] | ||
return attr, ok | ||
type _App struct { | ||
*_AppAgent | ||
|
||
isClose int32 | ||
} | ||
|
||
// NewApp 创建服务发现app | ||
func NewApp(checker IHealthChecker, container IAppContainer, attrs Attrs, nodes Nodes) IApp { | ||
if attrs == nil { | ||
attrs = make(Attrs) | ||
func (a *_App) Close() { | ||
if atomic.SwapInt32(&a.isClose, 1) == 0 { | ||
atomic.AddInt64(&a.use, -1) | ||
} | ||
return &app{ | ||
attrs: attrs, | ||
nodes: nodes, | ||
locker: sync.RWMutex{}, | ||
healthChecker: checker, | ||
id: uuid.NewString(), | ||
container: container, | ||
} | ||
} | ||
|
||
// ID 返回app的id | ||
func (s *app) ID() string { | ||
return s.id | ||
} | ||
|
||
// Nodes 将运行中的节点列表返回 | ||
func (s *app) Nodes() []INode { | ||
s.locker.RLock() | ||
defer s.locker.RUnlock() | ||
nodes := make([]INode, 0, len(s.nodes)) | ||
for _, node := range s.nodes { | ||
if node.Status() != Running { | ||
continue | ||
} | ||
nodes = append(nodes, node) | ||
} | ||
return nodes | ||
func newApp(nodes []eocontext.INode) *_AppAgent { | ||
|
||
return &_AppAgent{nodes: nodes} | ||
} | ||
|
||
// NodeError 定时检查节点,当节点失败时,则返回错误 | ||
func (s *app) NodeError(id string) error { | ||
s.locker.Lock() | ||
defer s.locker.Unlock() | ||
if n, ok := s.nodes[id]; ok { | ||
n.Down() | ||
if s.healthChecker != nil { | ||
err := s.healthChecker.AddToCheck(n) | ||
return err | ||
} | ||
} | ||
return nil | ||
func (a *_AppAgent) reset(nodes []eocontext.INode) { | ||
|
||
a.locker.Lock() | ||
defer a.locker.Unlock() | ||
a.nodes = nodes | ||
} | ||
|
||
// Close 关闭服务发现的app | ||
func (s *app) Close() error { | ||
// | ||
s.container.Remove(s.id) | ||
if s.healthChecker != nil { | ||
return s.healthChecker.Stop() | ||
} | ||
return nil | ||
func (a *_AppAgent) Nodes() []eocontext.INode { | ||
a.locker.RLock() | ||
defer a.locker.RUnlock() | ||
l := make([]eocontext.INode, len(a.nodes)) | ||
copy(l, a.nodes) | ||
return l | ||
} |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,8 @@ | ||
package discovery | ||
|
||
// IHealthChecker 健康检查接口 | ||
type IHealthChecker interface { | ||
Check(nodes INodes) | ||
Reset(conf interface{}) error | ||
Stop() | ||
} |
This file was deleted.
Oops, something went wrong.
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Oops, something went wrong.