1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
|
package registry
import (
"fmt"
"strconv"
"strings"
"github.com/go-chassis/go-chassis/core/registry"
"github.com/go-chassis/go-chassis/pkg/util/tags"
v1 "k8s.io/api/core/v1"
"k8s.io/klog/v2"
"github.com/kubeedge/kubeedge/edge/pkg/metamanager/client"
"github.com/kubeedge/kubeedge/edgemesh/pkg/cache"
"github.com/kubeedge/kubeedge/edgemesh/pkg/common"
)
const (
// EdgeRegistry constant string
EdgeRegistry = "edge"
)
// TODO Remove the init method, because it will cause invalid logs to be printed when the program is running @kadisi
// init initialize the plugin of edge meta registry
func init() { registry.InstallServiceDiscovery(EdgeRegistry, NewEdgeServiceDiscovery) }
// EdgeServiceDiscovery to represent the object of service center to call the APIs of service center
type EdgeServiceDiscovery struct {
metaClient client.CoreInterface
Name string
}
func NewEdgeServiceDiscovery(options registry.Options) registry.ServiceDiscovery {
return &EdgeServiceDiscovery{
metaClient: client.New(),
Name: EdgeRegistry,
}
}
// GetAllMicroServices Get all MicroService information.
func (esd *EdgeServiceDiscovery) GetAllMicroServices() ([]*registry.MicroService, error) {
return nil, nil
}
// FindMicroServiceInstances find micro-service instances (subnets)
func (esd *EdgeServiceDiscovery) FindMicroServiceInstances(consumerID, microServiceName string, tags utiltags.Tags) ([]*registry.MicroServiceInstance, error) {
// parse microServiceName
name, namespace, port, err := parseServiceURL(microServiceName)
if err != nil {
return nil, err
}
// get service
service, err := esd.getService(name, namespace)
if err != nil {
return nil, err
}
// get pods
pods, err := esd.getPods(name, namespace)
if err != nil {
return nil, err
}
// get targetPort
var targetPort int
for _, p := range service.Spec.Ports {
if p.Protocol == "TCP" && int(p.Port) == port {
targetPort = p.TargetPort.IntValue()
break
}
}
// port not found
if targetPort == 0 {
klog.Errorf("[EdgeMesh] port %d not found in svc: %s.%s", port, namespace, name)
return nil, fmt.Errorf("port %d not found in svc: %s.%s", port, namespace, name)
}
// gen
var microServiceInstances []*registry.MicroServiceInstance
var hostPort int32
// all pods share the same hostport, get from pods[0]
if pods[0].Spec.HostNetwork {
// host network
hostPort = int32(targetPort)
} else {
// container network
for _, container := range pods[0].Spec.Containers {
for _, port := range container.Ports {
if port.ContainerPort == int32(targetPort) {
hostPort = port.HostPort
}
}
}
}
for _, p := range pods {
if p.Status.Phase == v1.PodRunning {
microServiceInstances = append(microServiceInstances, ®istry.MicroServiceInstance{
InstanceID: "",
ServiceID: name + "." + namespace,
HostName: "",
EndpointsMap: map[string]string{"rest": fmt.Sprintf("%s:%d", p.Status.HostIP, hostPort)},
})
}
}
return microServiceInstances, nil
}
// GetMicroServiceID get microServiceID
func (esd *EdgeServiceDiscovery) GetMicroServiceID(appID, microServiceName, version, env string) (string, error) {
return "", nil
}
// GetMicroServiceInstances return instances
func (esd *EdgeServiceDiscovery) GetMicroServiceInstances(consumerID, providerID string) ([]*registry.MicroServiceInstance, error) {
return nil, nil
}
// GetMicroService return service
func (esd *EdgeServiceDiscovery) GetMicroService(microServiceID string) (*registry.MicroService, error) {
return nil, nil
}
// AutoSync updating the cache manager
func (esd *EdgeServiceDiscovery) AutoSync() {}
// Close close all websocket connection
func (esd *EdgeServiceDiscovery) Close() error { return nil }
// parseServiceURL parses serviceURL to ${service_name}.${namespace}.svc.${cluster}:${port}, keeps with k8s service
func parseServiceURL(serviceURL string) (string, string, int, error) {
var port int
var err error
serviceURLSplit := strings.Split(serviceURL, ":")
if len(serviceURLSplit) == 1 {
// default
port = 80
} else if len(serviceURLSplit) == 2 {
port, err = strconv.Atoi(serviceURLSplit[1])
if err != nil {
klog.Errorf("[EdgeMesh] service url %s invalid", serviceURL)
return "", "", 0, err
}
} else {
klog.Errorf("[EdgeMesh] service url %s invalid", serviceURL)
err = fmt.Errorf("service url %s invalid", serviceURL)
return "", "", 0, err
}
name, namespace := common.SplitServiceKey(serviceURLSplit[0])
return name, namespace, port, nil
}
// getService get k8s service from either lruCache or metaManager
func (esd *EdgeServiceDiscovery) getService(name, namespace string) (*v1.Service, error) {
var svc *v1.Service
key := fmt.Sprintf("service.%s.%s", namespace, name)
// try to get from cache
v, ok := cache.GetMeshCache().Get(key)
if ok {
svc, ok = v.(*v1.Service)
if !ok {
klog.Errorf("[EdgeMesh] service %s from cache with invalid type", key)
return nil, fmt.Errorf("service %s from cache with invalid type", key)
}
klog.Infof("[EdgeMesh] get service %s from cache", key)
} else {
// get from metaClient
var err error
svc, err = esd.metaClient.Services(namespace).Get(name)
if err != nil {
klog.Errorf("[EdgeMesh] get service from metaClient failed, error: %v", err)
return nil, err
}
cache.GetMeshCache().Add(key, svc)
klog.Infof("[EdgeMesh] get service %s from metaClient", key)
}
return svc, nil
}
// getPods get service pods from either lruCache or metaManager
func (esd *EdgeServiceDiscovery) getPods(name, namespace string) ([]v1.Pod, error) {
var pods []v1.Pod
key := fmt.Sprintf("pods.%s.%s", namespace, name)
// try to get from cache
v, ok := cache.GetMeshCache().Get(key)
if ok {
pods, ok = v.([]v1.Pod)
if !ok {
klog.Errorf("[EdgeMesh] pods %s from cache with invalid type", key)
return nil, fmt.Errorf("pods %s from cache with invalid type", key)
}
if len(pods) == 0 {
klog.Errorf("[EdgeMesh] pod list %s is empty", key)
return nil, fmt.Errorf("pod list %s is empty", key)
}
klog.Infof("[EdgeMesh] get pods %s from cache", key)
} else {
// get from metaClient
var err error
pods, err = esd.metaClient.Services(namespace).GetPods(name)
if err != nil {
klog.Errorf("[EdgeMesh] get pods from metaClient failed, error: %v", err)
return nil, err
}
if len(pods) == 0 {
klog.Errorf("[EdgeMesh] pod list %s is empty", key)
return nil, fmt.Errorf("pod list %s is empty", key)
}
cache.GetMeshCache().Add(key, pods)
klog.Infof("[EdgeMesh] get pods %s from metaClient", key)
}
return pods, nil
}
|