The current repo belongs to Paused status, and some functions are restricted. For details, please refer to the description of repo status
2 Star 0 Fork 1

JUMEI_ARCH/go-plugins
Paused

Create your Gitee Account
Explore and code with more than 12 million developers,Free private repositories !:)
Sign up
文件
Clone or Download
label.go 2.87 KB
Copy Edit Raw Blame History
Asim Aslam authored 2018-03-03 12:28 . switch to stdlib context
package label
import (
"context"
"sync"
"github.com/micro/go-micro/cmd"
"github.com/micro/go-micro/registry"
"github.com/micro/go-micro/selector"
)
type labelSelector struct {
so selector.Options
}
func init() {
cmd.DefaultSelectors["label"] = NewSelector
}
func prioritise(nodes []*registry.Node, labels []label) []*registry.Node {
var lnodes []*registry.Node
marked := make(map[string]bool)
for _, label := range labels {
for _, node := range nodes {
// already used
if _, ok := marked[node.Id]; ok {
continue
}
// nil metadata?
if node.Metadata == nil {
continue
}
// matching label?
if val, ok := node.Metadata[label.key]; !ok || label.val != val {
continue
}
// matched! mark it
marked[node.Id] = true
// append to nodes
lnodes = append(lnodes, node)
}
}
// grab the leftovers
for _, node := range nodes {
if _, ok := marked[node.Id]; ok {
continue
}
lnodes = append(lnodes, node)
}
return lnodes
}
func next(nodes []*registry.Node) func() (*registry.Node, error) {
var i int
var mtx sync.Mutex
return func() (*registry.Node, error) {
mtx.Lock()
if i >= len(nodes) {
i = 0
}
node := nodes[i]
i++
mtx.Unlock()
return node, nil
}
}
func (r *labelSelector) Init(opts ...selector.Option) error {
for _, o := range opts {
o(&r.so)
}
return nil
}
func (r *labelSelector) Options() selector.Options {
return r.so
}
func (r *labelSelector) Select(service string, opts ...selector.SelectOption) (selector.Next, error) {
var sopts selector.SelectOptions
for _, opt := range opts {
opt(&sopts)
}
// get the service
services, err := r.so.Registry.GetService(service)
if err != nil {
return nil, err
}
// apply the filters
for _, filter := range sopts.Filters {
services = filter(services)
}
// if there's nothing left, return
if len(services) == 0 {
return nil, selector.ErrNotFound
}
var nodes []*registry.Node
// flatten node list
for _, service := range services {
for _, node := range service.Nodes {
nodes = append(nodes, node)
}
}
// any nodes left?
if len(nodes) == 0 {
return nil, selector.ErrNotFound
}
// now prioritise the list based on labels
// oh god the O(n)^2 cruft or well not really
// more like O(m*n) or something like that
if labels, ok := r.so.Context.Value(labelKey{}).([]label); ok {
nodes = prioritise(nodes, labels)
}
return next(nodes), nil
}
func (r *labelSelector) Mark(service string, node *registry.Node, err error) {
return
}
func (r *labelSelector) Reset(service string) {
return
}
func (r *labelSelector) Close() error {
return nil
}
func (r *labelSelector) String() string {
return "label"
}
func NewSelector(opts ...selector.Option) selector.Selector {
sopts := selector.Options{
Context: context.TODO(),
Registry: registry.DefaultRegistry,
}
for _, opt := range opts {
opt(&sopts)
}
return &labelSelector{sopts}
}
马建仓 AI 助手
尝试更多
代码解读
代码找茬
代码优化
Go
1
https://gitee.com/JMArch/go-plugins.git
git@gitee.com:JMArch/go-plugins.git
JMArch
go-plugins
go-plugins
v0.9.3

Search

0d507c66 1850385 C8b1a773 1850385