summaryrefslogtreecommitdiff
path: root/edgemesh
diff options
context:
space:
mode:
authorliuzhiyi1993 <liuzhiyi@huawei.com>2020-02-13 05:56:21 +0800
committerliuzhiyi1993 <liuzhiyi@huawei.com>2020-03-31 22:30:02 +0800
commit0dfe63446ad76e41a084ade9d3c393a2554def34 (patch)
treef9857c1fd4ad3e9a00aa04effdb1c57bde6398b7 /edgemesh
parentMerge pull request #1567 from chendave/broken (diff)
downloadkubeedge-0dfe63446ad76e41a084ade9d3c393a2554def34.tar.gz
refactor edgemesh
Diffstat (limited to 'edgemesh')
-rw-r--r--edgemesh/cmd/edgemesh.go7
-rw-r--r--edgemesh/pkg/cache/cache.go20
-rw-r--r--edgemesh/pkg/common/utils.go16
-rw-r--r--edgemesh/pkg/common/utils_test.go29
-rw-r--r--edgemesh/pkg/config/config.go35
-rw-r--r--edgemesh/pkg/dns/dns.go436
-rw-r--r--edgemesh/pkg/listener/listener.go535
-rw-r--r--edgemesh/pkg/module.go25
-rw-r--r--edgemesh/pkg/plugin/panel/panel.go (renamed from edgemesh/pkg/panel/fake.go)17
-rw-r--r--edgemesh/pkg/plugin/plugin.go48
-rw-r--r--edgemesh/pkg/plugin/registry/registry.go204
-rw-r--r--edgemesh/pkg/protocol/http/http.go131
-rw-r--r--edgemesh/pkg/protocol/protocol.go5
-rw-r--r--edgemesh/pkg/proxy/poll/poll.go85
-rw-r--r--edgemesh/pkg/proxy/proxy.go657
-rw-r--r--edgemesh/pkg/proxy/virtualdevice/virtualdevice.go72
-rw-r--r--edgemesh/pkg/registry/registry.go95
-rw-r--r--edgemesh/pkg/resolver/resolver.go74
-rw-r--r--edgemesh/pkg/resolver/resolver_chain.go28
-rw-r--r--edgemesh/pkg/resolver/resolver_chain_test.go91
-rw-r--r--edgemesh/pkg/server/dns.go427
-rw-r--r--edgemesh/pkg/server/ip.go28
-rw-r--r--edgemesh/pkg/server/server.go49
-rw-r--r--edgemesh/pkg/server/tcp.go120
-rw-r--r--edgemesh/pkg/server/tcp_test.go249
25 files changed, 1628 insertions, 1855 deletions
diff --git a/edgemesh/cmd/edgemesh.go b/edgemesh/cmd/edgemesh.go
index c167f145e..9b1a3b665 100644
--- a/edgemesh/cmd/edgemesh.go
+++ b/edgemesh/cmd/edgemesh.go
@@ -2,12 +2,8 @@ package main
import (
"flag"
-
"github.com/spf13/pflag"
"k8s.io/klog"
-
- _ "github.com/kubeedge/kubeedge/edgemesh/pkg/panel"
- "github.com/kubeedge/kubeedge/edgemesh/pkg/server"
)
func main() {
@@ -17,7 +13,4 @@ func main() {
// TODO need parse edgemesh config file before Register @kadisi
//pkg.Register()
-
- //Start server
- server.StartTCP()
}
diff --git a/edgemesh/pkg/cache/cache.go b/edgemesh/pkg/cache/cache.go
new file mode 100644
index 000000000..63c31f6cb
--- /dev/null
+++ b/edgemesh/pkg/cache/cache.go
@@ -0,0 +1,20 @@
+package cache
+
+import "github.com/hashicorp/golang-lru"
+
+const DefaultCapacity = 20
+
+// meshCache is the lru cache to store edgemesh meta
+var meshCache *lru.Cache
+
+func init() {
+ var err error
+ meshCache, err = lru.New(DefaultCapacity)
+ if err != nil {
+ panic(err)
+ }
+}
+
+func GetMeshCache() *lru.Cache {
+ return meshCache
+}
diff --git a/edgemesh/pkg/common/utils.go b/edgemesh/pkg/common/utils.go
index 9366c7aa3..c664c9173 100644
--- a/edgemesh/pkg/common/utils.go
+++ b/edgemesh/pkg/common/utils.go
@@ -1,6 +1,8 @@
package common
import (
+ "fmt"
+ "net"
"os"
"strings"
@@ -22,3 +24,17 @@ func SplitServiceKey(key string) (name, namespace string) {
}
return key, ns
}
+
+func GetInterfaceIP(name string) (net.IP, error) {
+ ifi, err := net.InterfaceByName(name)
+ if err != nil {
+ return nil, err
+ }
+ addrs, _ := ifi.Addrs()
+ for _, addr := range addrs {
+ if ip, ipn, _ := net.ParseCIDR(addr.String()); len(ipn.Mask) == 4 {
+ return ip, nil
+ }
+ }
+ return nil, fmt.Errorf("no ip of version 4 found for interface %s", name)
+}
diff --git a/edgemesh/pkg/common/utils_test.go b/edgemesh/pkg/common/utils_test.go
deleted file mode 100644
index 4c371e352..000000000
--- a/edgemesh/pkg/common/utils_test.go
+++ /dev/null
@@ -1,29 +0,0 @@
-package common
-
-import "testing"
-
-func TestSplitServiceKey(t *testing.T) {
- type args struct {
- key string
- }
- tests := []struct {
- name string
- args args
- wantName string
- wantNamespace string
- }{
- {"normal serviceKey", args{"foo.bar"}, "foo", "bar"},
- {"default namespace", args{"foo"}, "foo", "default"},
- }
- for _, tt := range tests {
- t.Run(tt.name, func(t *testing.T) {
- gotName, gotNamespace := SplitServiceKey(tt.args.key)
- if gotName != tt.wantName {
- t.Errorf("SplitServiceKey() gotName = %v, want %v", gotName, tt.wantName)
- }
- if gotNamespace != tt.wantNamespace {
- t.Errorf("SplitServiceKey() gotNamespace = %v, want %v", gotNamespace, tt.wantNamespace)
- }
- })
- }
-}
diff --git a/edgemesh/pkg/config/config.go b/edgemesh/pkg/config/config.go
index 6267c678c..79bb37d4b 100644
--- a/edgemesh/pkg/config/config.go
+++ b/edgemesh/pkg/config/config.go
@@ -1,8 +1,12 @@
package config
import (
+ "net"
"sync"
+ "k8s.io/klog"
+
+ "github.com/kubeedge/kubeedge/edgemesh/pkg/common"
"github.com/kubeedge/kubeedge/pkg/apis/componentconfig/edgecore/v1alpha1"
)
@@ -11,6 +15,9 @@ var once sync.Once
type Configure struct {
v1alpha1.EdgeMesh
+ // for edgemesh listener
+ ListenIP net.IP
+ Listener *net.TCPListener
}
func InitConfigure(e *v1alpha1.EdgeMesh) {
@@ -18,5 +25,33 @@ func InitConfigure(e *v1alpha1.EdgeMesh) {
Config = Configure{
EdgeMesh: *e,
}
+ if Config.Enable {
+ // get listen ip
+ var err error
+ Config.ListenIP, err = common.GetInterfaceIP(Config.ListenInterface)
+ if err != nil {
+ klog.Errorf("[EdgeMesh] get listen ip err: %v", err)
+ return
+ }
+ // get listener
+ tmpPort := 0
+ listenAddr := &net.TCPAddr{
+ IP: Config.ListenIP,
+ Port: Config.ListenPort + tmpPort,
+ }
+ for {
+ ln, err := net.ListenTCP("tcp", listenAddr)
+ if err == nil {
+ Config.Listener = ln
+ break
+ }
+ klog.Warningf("[EdgeMesh] listen on address %v err: %v", listenAddr, err)
+ tmpPort++
+ listenAddr = &net.TCPAddr{
+ IP: Config.ListenIP,
+ Port: Config.ListenPort + tmpPort,
+ }
+ }
+ }
})
}
diff --git a/edgemesh/pkg/dns/dns.go b/edgemesh/pkg/dns/dns.go
new file mode 100644
index 000000000..3905f9bc0
--- /dev/null
+++ b/edgemesh/pkg/dns/dns.go
@@ -0,0 +1,436 @@
+package dns
+
+import (
+ "bufio"
+ "encoding/binary"
+ "errors"
+ "fmt"
+ "net"
+ "os"
+ "strings"
+ "time"
+ "unsafe"
+
+ "k8s.io/klog"
+
+ "github.com/kubeedge/kubeedge/edge/pkg/metamanager/client"
+ "github.com/kubeedge/kubeedge/edgemesh/pkg/common"
+ "github.com/kubeedge/kubeedge/edgemesh/pkg/config"
+ "github.com/kubeedge/kubeedge/edgemesh/pkg/listener"
+)
+
+type Event int
+
+var (
+ // default docker0
+ ifi = "docker0"
+ // QR: 0 represents query, 1 represents response
+ dnsQR = uint16(0x8000)
+ oneByteSize = uint16(1)
+ twoByteSize = uint16(2)
+ ttl = uint32(64)
+)
+
+const (
+ // 1 for ipv4
+ aRecord = 1
+ bufSize = 1024
+ errNotImplemented = uint16(0x0004)
+ errRefused = uint16(0x0005)
+ eventNothing = Event(0)
+ eventUpstream = Event(1)
+ eventNxDomain = Event(2)
+)
+
+type dnsHeader struct {
+ id uint16
+ flags uint16
+ qdCount uint16
+ anCount uint16
+ nsCount uint16
+ arCount uint16
+}
+
+type dnsQuestion struct {
+ from *net.UDPAddr
+ head *dnsHeader
+ name []byte
+ queByte []byte
+ qType uint16
+ qClass uint16
+ queNum uint16
+ event Event
+}
+
+type dnsAnswer struct {
+ name []byte
+ qType uint16
+ qClass uint16
+ ttl uint32
+ dataLen uint16
+ addr []byte
+}
+
+// metaClient is a query client
+var metaClient client.CoreInterface
+
+// dnsConn saves DNS protocol
+var dnsConn *net.UDPConn
+
+// Start is for external call
+func Start() {
+ startDNS()
+}
+
+// startDNS starts edgemesh dns server
+func startDNS() {
+ // init meta client
+ metaClient = client.New()
+ // get dns listen ip
+ lip, err := common.GetInterfaceIP(ifi)
+ if err != nil {
+ klog.Errorf("[EdgeMesh] get dns listen ip err: %v", err)
+ return
+ }
+
+ laddr := &net.UDPAddr{
+ IP: lip,
+ Port: 53,
+ }
+ udpConn, err := net.ListenUDP("udp", laddr)
+ if err != nil {
+ klog.Errorf("[EdgeMesh] dns server listen on %v error: %v", laddr, err)
+ return
+ }
+ defer udpConn.Close()
+ dnsConn = udpConn
+ for {
+ req := make([]byte, bufSize)
+ n, from, err := dnsConn.ReadFromUDP(req)
+ if err != nil || n <= 0 {
+ klog.Errorf("[EdgeMesh] dns server read from udp error: %v", err)
+ continue
+ }
+
+ que, err := parseDNSQuery(req[:n])
+ if err != nil {
+ continue
+ }
+
+ que.from = from
+
+ rsp := make([]byte, 0)
+ rsp, err = recordHandle(que, req[:n])
+ if err != nil {
+ klog.Warningf("[EdgeMesh] failed to resolve dns: %v", err)
+ continue
+ }
+ dnsConn.WriteTo(rsp, from)
+ }
+}
+
+// recordHandle returns the answer for the dns question
+func recordHandle(que *dnsQuestion, req []byte) (rsp []byte, err error) {
+ var exist bool
+ var ip string
+ // qType should be 1 for ipv4
+ if que.name != nil && que.qType == aRecord {
+ domainName := string(que.name)
+ exist, ip = lookupFromMetaManager(domainName)
+ }
+
+ if !exist || que.event == eventUpstream {
+ // if this service doesn't belongs to this cluster
+ go getFromRealDNS(req, que.from)
+ return rsp, fmt.Errorf("get from real dns")
+ }
+
+ address := net.ParseIP(ip).To4()
+ if address == nil {
+ que.event = eventNxDomain
+ }
+ // gen
+ pre := modifyRspPrefix(que)
+ rsp = append(rsp, pre...)
+ if que.event != eventNothing {
+ return rsp, nil
+ }
+ // create a deceptive resp, if no error
+ dnsAns := &dnsAnswer{
+ name: que.name,
+ qType: que.qType,
+ qClass: que.qClass,
+ ttl: ttl,
+ dataLen: uint16(len(address)),
+ addr: address,
+ }
+ ans := dnsAns.getAnswer()
+ rsp = append(rsp, ans...)
+
+ return rsp, nil
+}
+
+// parseDNSQuery converts bytes to *dnsQuestion
+func parseDNSQuery(req []byte) (que *dnsQuestion, err error) {
+ head := &dnsHeader{}
+ head.getHeader(req)
+ if !head.isAQuery() {
+ return nil, errors.New("not a dns query, ignore")
+ }
+ que = &dnsQuestion{
+ event: eventNothing,
+ }
+ // Generally, when the recursive DNS server requests upward, it may
+ // initiate a resolution request for multiple aliases/domain names
+ // at once, Edge DNS does not need to process a message that carries
+ // multiple questions at a time.
+ if head.qdCount != 1 {
+ que.event = eventUpstream
+ return
+ }
+
+ offset := uint16(unsafe.Sizeof(dnsHeader{}))
+ // DNS NS <ROOT> operation
+ if req[offset] == 0x0 {
+ que.event = eventUpstream
+ return
+ }
+ que.getQuestion(req, offset, head)
+ err = nil
+ return
+}
+
+// isAQuery judges if the dns pkg is a query
+func (h *dnsHeader) isAQuery() bool {
+ if h.flags&dnsQR != dnsQR {
+ return true
+ }
+ return false
+}
+
+// getHeader gets dns pkg head
+func (h *dnsHeader) getHeader(req []byte) {
+ h.id = binary.BigEndian.Uint16(req[0:2])
+ h.flags = binary.BigEndian.Uint16(req[2:4])
+ h.qdCount = binary.BigEndian.Uint16(req[4:6])
+ h.anCount = binary.BigEndian.Uint16(req[6:8])
+ h.nsCount = binary.BigEndian.Uint16(req[8:10])
+ h.arCount = binary.BigEndian.Uint16(req[10:12])
+}
+
+// getQuestion gets a dns question
+func (q *dnsQuestion) getQuestion(req []byte, offset uint16, head *dnsHeader) {
+ ost := offset
+ tmp := ost
+ ost = q.getQName(req, ost)
+ q.qType = binary.BigEndian.Uint16(req[ost : ost+twoByteSize])
+ ost += twoByteSize
+ q.qClass = binary.BigEndian.Uint16(req[ost : ost+twoByteSize])
+ ost += twoByteSize
+ q.head = head
+ q.queByte = req[tmp:ost]
+}
+
+// getAnswer generates answer for the dns question
+func (da *dnsAnswer) getAnswer() (answer []byte) {
+ answer = make([]byte, 0)
+
+ if da.qType == aRecord {
+ answer = append(answer, 0xc0)
+ answer = append(answer, 0x0c)
+
+ tmp16 := make([]byte, 2)
+ tmp32 := make([]byte, 4)
+
+ binary.BigEndian.PutUint16(tmp16, da.qType)
+ answer = append(answer, tmp16...)
+ binary.BigEndian.PutUint16(tmp16, da.qClass)
+ answer = append(answer, tmp16...)
+ binary.BigEndian.PutUint32(tmp32, da.ttl)
+ answer = append(answer, tmp32...)
+ binary.BigEndian.PutUint16(tmp16, da.dataLen)
+ answer = append(answer, tmp16...)
+ answer = append(answer, da.addr...)
+ }
+
+ return answer
+}
+
+// getQName gets dns question qName
+func (q *dnsQuestion) getQName(req []byte, offset uint16) uint16 {
+ ost := offset
+
+ for {
+ // one byte to suggest length
+ qbyte := uint16(req[ost])
+
+ // qName ends with 0x00, and 0x00 should not be included
+ if qbyte == 0x00 {
+ q.name = q.name[:uint16(len(q.name))-oneByteSize]
+ return ost + oneByteSize
+ }
+ // step forward one more byte and get the real stuff
+ ost += oneByteSize
+ q.name = append(q.name, req[ost:ost+qbyte]...)
+ // add "." symbol
+ q.name = append(q.name, 0x2e)
+ ost += qbyte
+ }
+}
+
+// lookupFromMetaManager confirms if the service exists
+func lookupFromMetaManager(serviceUrl string) (exist bool, ip string) {
+ name, namespace := common.SplitServiceKey(serviceUrl)
+ s, _ := metaClient.Services(namespace).Get(name)
+ if s != nil {
+ svcName := namespace + "." + name
+ ip := listener.GetServiceServer(svcName)
+ klog.Infof("[EdgeMesh] dns server parse %s ip %s", serviceUrl, ip)
+ return true, ip
+ }
+ klog.Errorf("[EdgeMesh] service %s is not found in this cluster", serviceUrl)
+ return false, ""
+}
+
+// getFromRealDNS returns a dns response from real dns servers
+func getFromRealDNS(req []byte, from *net.UDPAddr) {
+ rsp := make([]byte, 0)
+ ips, err := parseNameServer()
+ if err != nil {
+ klog.Errorf("[EdgeMesh] parse nameserver err: %v", err)
+ return
+ }
+
+ laddr := &net.UDPAddr{
+ IP: net.IPv4zero,
+ Port: 0,
+ }
+
+ // get from real dns servers
+ for _, ip := range ips {
+ raddr := &net.UDPAddr{
+ IP: ip,
+ Port: 53,
+ }
+ conn, err := net.DialUDP("udp", laddr, raddr)
+ if err != nil {
+ continue
+ }
+ defer conn.Close()
+ _, err = conn.Write(req)
+ if err != nil {
+ continue
+ }
+ if err = conn.SetReadDeadline(time.Now().Add(time.Minute)); err != nil {
+ continue
+ }
+ var n int
+ buf := make([]byte, bufSize)
+ n, err = conn.Read(buf)
+ if err != nil {
+ continue
+ }
+
+ if n > 0 {
+ rsp = append(rsp, buf[:n]...)
+ dnsConn.WriteToUDP(rsp, from)
+ break
+ }
+ }
+}
+
+// parseNameServer gets all real nameservers from the resolv.conf
+func parseNameServer() ([]net.IP, error) {
+ file, err := os.Open("/etc/resolv.conf")
+ if err != nil {
+ return nil, fmt.Errorf("error opening /etc/resolv.conf: %v", err)
+ }
+ defer file.Close()
+
+ scan := bufio.NewScanner(file)
+ scan.Split(bufio.ScanLines)
+
+ ip := make([]net.IP, 0)
+
+ for scan.Scan() {
+ serverString := scan.Text()
+ if strings.Contains(serverString, "nameserver") {
+ tmpString := strings.Replace(serverString, "nameserver", "", 1)
+ nameserver := strings.TrimSpace(tmpString)
+ sip := net.ParseIP(nameserver)
+ if sip != nil && !sip.Equal(config.Config.ListenIP) {
+ ip = append(ip, sip)
+ }
+ }
+ }
+ if len(ip) == 0 {
+ return nil, fmt.Errorf("there is no nameserver in /etc/resolv.conf")
+ }
+ return ip, nil
+}
+
+// modifyRspPrefix generates a dns response head
+func modifyRspPrefix(que *dnsQuestion) (pre []byte) {
+ if que == nil {
+ return
+ }
+ // use head in que
+ rspHead := que.head
+ rspHead.convertQueryRsp(true)
+ if que.qType == aRecord {
+ rspHead.setAnswerNum(1)
+ } else {
+ rspHead.setAnswerNum(0)
+ }
+
+ rspHead.setRspRCode(que)
+ pre = rspHead.getByteFromDNSHeader()
+
+ pre = append(pre, que.queByte...)
+ return
+}
+
+// convertQueryRsp converts a dns question head to a response head
+func (h *dnsHeader) convertQueryRsp(isRsp bool) {
+ if isRsp {
+ h.flags |= dnsQR
+ } else {
+ h.flags |= dnsQR
+ }
+}
+
+// setAnswerNum sets the answer num for dns head
+func (h *dnsHeader) setAnswerNum(num uint16) {
+ h.anCount = num
+}
+
+// setRspRCode sets dns response return code
+func (h *dnsHeader) setRspRCode(que *dnsQuestion) {
+ if que.qType != aRecord {
+ h.flags &= (^errNotImplemented)
+ h.flags |= errNotImplemented
+ } else if que.event == eventNxDomain {
+ h.flags &= (^errRefused)
+ h.flags |= errRefused
+ }
+}
+
+// getByteFromDNSHeader converts dnsHeader to bytes
+func (h *dnsHeader) getByteFromDNSHeader() (rspHead []byte) {
+ rspHead = make([]byte, unsafe.Sizeof(*h))
+
+ idxTransactionID := unsafe.Sizeof(h.id)
+ idxFlags := unsafe.Sizeof(h.flags) + idxTransactionID
+ idxQDCount := unsafe.Sizeof(h.anCount) + idxFlags
+ idxANCount := unsafe.Sizeof(h.anCount) + idxQDCount
+ idxNSCount := unsafe.Sizeof(h.nsCount) + idxANCount
+ idxARCount := unsafe.Sizeof(h.arCount) + idxNSCount
+
+ binary.BigEndian.PutUint16(rspHead[:idxTransactionID], h.id)
+ binary.BigEndian.PutUint16(rspHead[idxTransactionID:idxFlags], h.flags)
+ binary.BigEndian.PutUint16(rspHead[idxFlags:idxQDCount], h.qdCount)
+ binary.BigEndian.PutUint16(rspHead[idxQDCount:idxANCount], h.anCount)
+ binary.BigEndian.PutUint16(rspHead[idxANCount:idxNSCount], h.nsCount)
+ binary.BigEndian.PutUint16(rspHead[idxNSCount:idxARCount], h.arCount)
+ return
+}
diff --git a/edgemesh/pkg/listener/listener.go b/edgemesh/pkg/listener/listener.go
new file mode 100644
index 000000000..e837296a7
--- /dev/null
+++ b/edgemesh/pkg/listener/listener.go
@@ -0,0 +1,535 @@
+package listener
+
+import (
+ "encoding/json"
+ "fmt"
+ "net"
+ "strconv"
+ "strings"
+ "sync"
+ "syscall"
+ "unsafe"
+
+ "github.com/kubeedge/beehive/pkg/core/model"
+ v1 "k8s.io/api/core/v1"
+ "k8s.io/klog"
+
+ "github.com/kubeedge/kubeedge/common/constants"
+ "github.com/kubeedge/kubeedge/edge/pkg/metamanager/client"
+ "github.com/kubeedge/kubeedge/edgemesh/pkg/cache"
+ "github.com/kubeedge/kubeedge/edgemesh/pkg/config"
+ "github.com/kubeedge/kubeedge/edgemesh/pkg/protocol"
+ "github.com/kubeedge/kubeedge/edgemesh/pkg/protocol/http"
+)
+
+type SvcDescription struct {
+ sync.RWMutex
+ SvcPortsByIP map[string]string // key: fakeIP, value: SvcPorts
+ IPBySvc map[string]string // key: svcName.svcNamespace, value: fakeIP
+}
+
+type sockAddr struct {
+ family uint16
+ data [14]byte
+}
+
+const (
+ defaultNetworkPrefix = "9.251."
+ maxPoolSize = 65534
+
+ SO_ORIGINAL_DST = 80
+)
+
+var (
+ svcDesc *SvcDescription
+ unused []string
+ indexOfPool uint16
+ metaClient client.CoreInterface
+ once sync.Once
+)
+
+func Init() {
+ once.Do(func() {
+ unused = make([]string, 0)
+ svcDesc = &SvcDescription{
+ SvcPortsByIP: make(map[string]string),
+ IPBySvc: make(map[string]string),
+ }
+ // init meta client
+ metaClient = client.New()
+ // init fakeIP pool
+ initPool()
+ // recover listener meta from edge db
+ recoverFromDB()
+ })
+}
+
+// getSubNet converts uint16 to "uint8.uint8"
+func getSubNet(subNet uint16) string {
+ arg1 := uint64(subNet & 0x00ff)
+ arg2 := uint64((subNet & 0xff00) >> 8)
+ return strconv.FormatUint(arg2, 10) + "." + strconv.FormatUint(arg1, 10)
+}
+
+// initPool initializes fakeIP pool with size of 256
+func initPool() {
+ // avoid 0.0
+ indexOfPool = uint16(1)
+ for ; indexOfPool <= uint16(255); indexOfPool++ {
+ ip := defaultNetworkPrefix + getSubNet(indexOfPool)
+ unused = append(unused, ip)
+ }
+}
+
+// expandPool expands fakeIP pool, each time with size of 256
+func expandPool() {
+ end := indexOfPool + uint16(255)
+ for ; indexOfPool <= end; indexOfPool++ {
+ // avoid 255.255
+ if indexOfPool > maxPoolSize {
+ return
+ }
+ ip := defaultNetworkPrefix + getSubNet(indexOfPool)
+ // if ip is not used, append it to unused
+ if svcDesc.getSvcPorts(ip) == "" {
+ unused = append(unused, ip)
+ }
+ }
+}
+
+// reserveIp reserves used fakeIP
+func reserveIP(ip string) {
+ for i, value := range unused {
+ if ip == value {
+ unused = append(unused[:i], unused[i+1:]...)
+ break
+ }
+ }
+}
+
+// recoverFromDB gets fakeIP from edge db and assigns them to services after EdgeMesh starts
+func recoverFromDB() {
+ svcs, err := metaClient.Services("all").ListAll()
+ if err != nil {
+ klog.Errorf("[EdgeMesh] list all services from edge db error: %v", err)
+ return
+ }
+ for _, svc := range svcs {
+ svcName := svc.Namespace + "." + svc.Name
+ value, err := metaClient.Listener().Get(svcName)
+ if err != nil {
+ klog.Errorf("[EdgeMesh] get listener of svc %s from edge db error: %v", svcName, err)
+ continue
+ }
+ ip, ok := value.([]string)
+ if !ok {
+ klog.Errorf("[EdgeMesh] value %+v is not a string", value)
+ continue
+ }
+ if len(ip) == 0 {
+ svcPorts := getSvcPorts(svc, svcName)
+ addServer(svcName, svcPorts)
+ klog.Warningf("[EdgeMesh] listener %s from edge db with no ip", svcName)
+ continue
+ }
+ svcPorts := getSvcPorts(svc, svcName)
+ reserveIP(ip[0][1 : len(ip[0])-1])
+ svcDesc.set(svcName, ip[0][1:len(ip[0])-1], svcPorts)
+ klog.Infof("[EdgeMesh] get listener %s from edge db: %s", svcName, ip[0][1:len(ip[0])-1])
+ }
+}
+
+// Start starts the EdgeMesh listener
+func Start() {
+ for {
+ conn, err := config.Config.Listener.Accept()
+ if err != nil {
+ klog.Warningf("[EdgeMesh] get tcp conn error: %v", err)
+ continue
+ }
+ ip, port, err := realServerAddress(&conn)
+ if err != nil {
+ klog.Warningf("[EdgeMesh] get real destination of tcp conn error: %v", err)
+ conn.Close()
+ continue
+ }
+ proto, err := newProtocolFromSock(ip, port, conn)
+ if err != nil {
+ klog.Warningf("[EdgeMesh] get protocol from sock err: %v", err)
+ conn.Close()
+ continue
+ }
+
+ go proto.Process()
+ }
+}
+
+// newProtocolFromSock returns a protocol.Protocol interface if the ip is in proxy list
+func newProtocolFromSock(ip string, port int, conn net.Conn) (proto protocol.Protocol, err error) {
+ svcPorts := svcDesc.getSvcPorts(ip)
+ protoName, svcName := getProtocol(svcPorts, port)
+ if protoName == "" || svcName == "" {
+ return nil, fmt.Errorf("protocol name: %s or svcName: %s is invalid", protoName, svcName)
+ }
+
+ svcNameSets := strings.Split(svcName, ".")
+ if len(svcNameSets) != 2 {
+ return nil, fmt.Errorf("invalid length %d after splitting svc name %s", len(svcNameSets), svcName)
+ }
+ namespace := svcNameSets[0]
+ name := svcNameSets[1]
+
+ switch protoName {
+ case "http":
+ proto = &http.HTTP{
+ Conn: conn,
+ SvcName: name,
+ SvcNamespace: namespace,
+ Port: port,
+ }
+ err = nil
+ default:
+ proto = nil
+ err = fmt.Errorf("protocol: %s is not supported yet", protoName)
+ }
+ return
+}
+
+// getProtocol gets protocol name
+func getProtocol(svcPorts string, port int) (string, string) {
+ var protoName string
+ sub := strings.Split(svcPorts, "|")
+ n := len(sub)
+ if n < 2 {
+ return "", ""
+ }
+ svcName := sub[n-1]
+
+ pstr := strconv.Itoa(port)
+ if pstr == "" {
+ return "", ""
+ }
+ for _, s := range sub {
+ if strings.Contains(s, pstr) {
+ protoName = strings.Split(s, ",")[0]
+ break
+ }
+ }
+ return protoName, svcName
+}
+
+// realServerAddress returns an intercepted connection's original destination.
+func realServerAddress(conn *net.Conn) (string, int, error) {
+ tcpConn, ok := (*conn).(*net.TCPConn)
+ if !ok {
+ return "", -1, fmt.Errorf("not a TCPConn")
+ }
+
+ file, err := tcpConn.File()
+ if err != nil {
+ return "", -1, err
+ }
+
+ // To avoid potential problems from making the socket non-blocking.
+ tcpConn.Close()
+ *conn, err = net.FileConn(file)
+ if err != nil {
+ return "", -1, err
+ }
+
+ defer file.Close()
+ fd := file.Fd()
+
+ var addr sockAddr
+ size := uint32(unsafe.Sizeof(addr))
+ err = getSockOpt(int(fd), syscall.SOL_IP, SO_ORIGINAL_DST, uintptr(unsafe.Pointer(&addr)), &size)
+ if err != nil {
+ return "", -1, err
+ }
+
+ var ip net.IP
+ switch addr.family {
+ case syscall.AF_INET:
+ ip = addr.data[2:6]
+ default:
+ return "", -1, fmt.Errorf("unrecognized address family")
+ }
+
+ port := int(addr.data[0])<<8 + int(addr.data[1])
+ syscall.SetNonblock(int(fd), true)
+
+ return ip.String(), port, nil
+}
+
+func getSockOpt(s int, level int, name int, val uintptr, vallen *uint32) (err error) {
+ _, _, e1 := syscall.Syscall6(syscall.SYS_GETSOCKOPT, uintptr(s), uintptr(level), uintptr(name), uintptr(val), uintptr(unsafe.Pointer(vallen)), 0)
+ if e1 != 0 {
+ err = e1
+ }
+ return
+}
+
+// filterResourceTypeService implements msg filter for "Service" and "ServiceList" resource
+func filterResourceTypeService(msg model.Message) []v1.Service {
+ svcs := make([]v1.Service, 0)
+ content, err := json.Marshal(msg.GetContent())
+ if err != nil || len(content) == 0 {
+ return svcs
+ }
+ switch getResourceType(msg.GetResource()) {
+ case constants.ResourceTypeService:
+ s, err := handleServiceMessage(content)
+ if err != nil {
+ break
+ }
+ svcs = append(svcs, *s)
+ case constants.ResourceTypeServiceList:
+ ss, err := handleServiceListMessage(content)
+ if err != nil {
+ break
+ }
+ svcs = append(svcs, ss...)
+ default:
+ break
+ }
+
+ return svcs
+}
+
+func getSvcPorts(svc v1.Service, svcName string) string {
+ svcPorts := ""
+ for _, p := range svc.Spec.Ports {
+ pro := strings.Split(p.Name, "-")
+ sub := fmt.Sprintf("%s,%d,%d|", pro[0], p.Port, p.TargetPort.IntVal)
+ svcPorts = svcPorts + sub
+ }
+ svcPorts += svcName
+ return svcPorts
+}
+
+// MsgProcess processes messages from metaManager
+func MsgProcess(msg model.Message) {
+ // process services
+ if svcs := filterResourceTypeService(msg); len(svcs) != 0 {
+ klog.Infof("[EdgeMesh] %s services: %d resource: %s", msg.GetOperation(), len(svcs), msg.Router.Resource)
+ for _, svc := range svcs {
+ svcName := svc.Namespace + "." + svc.Name
+ svcPorts := getSvcPorts(svc, svcName)
+ switch msg.GetOperation() {
+ case "insert":
+ cache.GetMeshCache().Add("service"+"."+svcName, &svc)
+ klog.Infof("[EdgeMesh] insert svc %s.%s into cache", svc.Namespace, svc.Name)
+ addServer(svcName, svcPorts)
+ case "update":
+ cache.GetMeshCache().Add("service"+"."+svcName, &svc)
+ klog.Infof("[EdgeMesh] update svc %s.%s in cache", svc.Namespace, svc.Name)
+ updateServer(svcName, svcPorts)
+ case "delete":
+ cache.GetMeshCache().Remove("service" + "." + svcName)
+ klog.Infof("[EdgeMesh] delete svc %s.%s from cache", svc.Namespace, svc.Name)
+ delServer(svcName)
+ default:
+ klog.Warningf("[EdgeMesh] invalid %s operation on services", msg.GetOperation())
+ }
+ }
+ return
+ }
+ // process pods
+ if getResourceType(msg.GetResource()) == model.ResourceTypePodlist {
+ klog.Infof("[EdgeMesh] %s podlist, resource: %s", msg.GetOperation(), msg.Router.Resource)
+ pods := make([]v1.Pod, 0)
+ content, err := json.Marshal(msg.GetContent())
+ if err != nil {
+ klog.Errorf("[EdgeMesh] marshal podlist msg content err: %v", err)
+ return
+ }
+ pods, err = handlePodListMessage(content)
+ if err != nil {
+ return
+ }
+ podListName := getResourceName(msg.GetResource())
+ podListNamespace := getResourceNamespace(msg.GetResource())
+ switch msg.GetOperation() {
+ case "insert", "update":
+ cache.GetMeshCache().Add("pods"+"."+podListNamespace+"."+podListName, pods)
+ klog.Infof("[EdgeMesh] insert/update pods %s.%s into cache", podListNamespace, podListName)
+ case "delete":
+ cache.GetMeshCache().Remove("pods" + "." + podListNamespace + "." + podListName)
+ klog.Infof("[EdgeMesh] delete pods %s.%s from cache", podListNamespace, podListName)
+ default:
+ klog.Warningf("[EdgeMesh] invalid %s operation on podlist", msg.GetOperation())
+ }
+ }
+ return
+}
+
+// addServer adds a server
+func addServer(svcName, svcPorts string) {
+ ip := svcDesc.getIP(svcName)
+ if ip != "" {
+ svcDesc.set(svcName, ip, svcPorts)
+ return
+ } else {
+ if len(unused) == 0 {
+ // try to expand
+ expandPool()
+ if len(unused) == 0 {
+ klog.Warningf("[EdgeMesh] insufficient fake IP !!")
+ return
+ }
+ }
+ ip = unused[0]
+ unused = unused[1:]
+ }
+ svcDesc.set(svcName, ip, svcPorts)
+ err := metaClient.Listener().Add(svcName, ip)
+ if err != nil {
+ klog.Errorf("[EdgeMesh] add listener %s to edge db error: %v", svcName, err)
+ return
+ }
+ return
+}
+
+// updateServer updates a server
+func updateServer(svcName, svcPorts string) {
+ ip := svcDesc.getIP(svcName)
+ if ip == "" {
+ if len(unused) == 0 {
+ // try to expand
+ expandPool()
+ if len(unused) == 0 {
+ klog.Warningf("[EdgeMesh] insufficient fake IP !!")
+ return
+ }
+ }
+ ip = unused[0]
+ unused = unused[1:]
+ err := metaClient.Listener().Add(svcName, ip)
+ if err != nil {
+ klog.Errorf("[EdgeMesh] add listener %s to edge db error: %v", svcName, err)
+ }
+ }
+ svcDesc.set(svcName, ip, svcPorts)
+}
+
+// delServer deletes a server
+func delServer(svcName string) {
+ ip := svcDesc.getIP(svcName)
+ if ip == "" {
+ return
+ }
+ svcDesc.del(svcName, ip)
+ err := metaClient.Listener().Del(svcName)
+ if err != nil {
+ klog.Errorf("[EdgeMesh] delete listener from edge db error: %v", err)
+ }
+ // recycling fakeIP
+ unused = append(unused, ip)
+}
+
+// handleServiceMessage converts bytes to k8s service meta
+func handleServiceMessage(content []byte) (*v1.Service, error) {
+ var s v1.Service
+ err := json.Unmarshal(content, &s)
+ if err != nil {
+ klog.Errorf("[EdgeMesh] unmarshal message to k8s service failed, err: %v", err)
+ return nil, err
+ }
+ return &s, nil
+}
+
+// handleServiceListMessage converts bytes to k8s service list meta
+func handleServiceListMessage(content []byte) ([]v1.Service, error) {
+ var ss []v1.Service
+ err := json.Unmarshal(content, &ss)
+ if err != nil {
+ klog.Errorf("[EdgeMesh] unmarshal message to k8s service list failed, err: %v", err)
+ return nil, err
+ }
+ return ss, nil
+}
+
+// handlePodListMessage converts bytes to k8s pod list meta
+func handlePodListMessage(content []byte) ([]v1.Pod, error) {
+ var pp []v1.Pod
+ err := json.Unmarshal(content, &pp)
+ if err != nil {
+ klog.Errorf("[EdgeMesh] unmarshal message to k8s pod list failed, err: %v", err)
+ return nil, err
+ }
+ return pp, nil
+}
+
+// getResourceType returns the resource type as a string
+func getResourceType(resource string) string {
+ str := strings.Split(resource, "/")
+ if len(str) == 3 {
+ return str[1]
+ } else if len(str) == 5 {
+ return str[3]
+ } else {
+ return resource
+ }
+}
+
+// getResourceName returns the resource name as a string
+func getResourceName(resource string) string {
+ str := strings.Split(resource, "/")
+ if len(str) == 3 {
+ return str[2]
+ } else if len(str) == 5 {
+ return str[4]
+ } else {
+ return resource
+ }
+}
+
+// getResourceNamespace returns the resource namespace as a string
+func getResourceNamespace(resource string) string {
+ str := strings.Split(resource, "/")
+ if len(str) == 3 {
+ return str[0]
+ } else if len(str) == 5 {
+ return str[2]
+ } else {
+ return resource
+ }
+}
+
+// set is a thread-safe operation to add to map
+func (sd *SvcDescription) set(svcName, ip, svcPorts string) {
+ sd.Lock()
+ defer sd.Unlock()
+ sd.IPBySvc[svcName] = ip
+ sd.SvcPortsByIP[ip] = svcPorts
+}
+
+// del is a thread-safe operation to del from map
+func (sd *SvcDescription) del(svcName, ip string) {
+ sd.Lock()
+ defer sd.Unlock()
+ delete(sd.IPBySvc, svcName)
+ delete(sd.SvcPortsByIP, ip)
+}
+
+// getIP is a thread-safe operation to get from map
+func (sd *SvcDescription) getIP(svcName string) string {
+ sd.RLock()
+ defer sd.RUnlock()
+ ip := sd.IPBySvc[svcName]
+ return ip
+}
+
+// getSvcPorts is a thread-safe operation to get from map
+func (sd *SvcDescription) getSvcPorts(ip string) string {
+ sd.RLock()
+ defer sd.RUnlock()
+ svcPorts := sd.SvcPortsByIP[ip]
+ return svcPorts
+}
+
+// GetServiceServer returns the proxy IP by given service name
+func GetServiceServer(svcName string) string {
+ ip := svcDesc.getIP(svcName)
+ return ip
+}
diff --git a/edgemesh/pkg/module.go b/edgemesh/pkg/module.go
index 285f7c0e5..ff4e58a08 100644
--- a/edgemesh/pkg/module.go
+++ b/edgemesh/pkg/module.go
@@ -1,15 +1,17 @@
package pkg
import (
- "k8s.io/klog"
-
"github.com/kubeedge/beehive/pkg/core"
beehiveContext "github.com/kubeedge/beehive/pkg/core/context"
+ "k8s.io/klog"
+
"github.com/kubeedge/kubeedge/edge/pkg/common/modules"
meshconfig "github.com/kubeedge/kubeedge/edgemesh/pkg/config"
"github.com/kubeedge/kubeedge/edgemesh/pkg/constant"
+ "github.com/kubeedge/kubeedge/edgemesh/pkg/dns"
+ "github.com/kubeedge/kubeedge/edgemesh/pkg/listener"
+ "github.com/kubeedge/kubeedge/edgemesh/pkg/plugin"
"github.com/kubeedge/kubeedge/edgemesh/pkg/proxy"
- "github.com/kubeedge/kubeedge/edgemesh/pkg/server"
"github.com/kubeedge/kubeedge/pkg/apis/componentconfig/edgecore/v1alpha1"
)
@@ -41,22 +43,31 @@ func (em *EdgeMesh) Enable() bool {
//Start sets context and starts the controller
func (em *EdgeMesh) Start() {
+ // install go-chassis plugins
+ plugin.Install()
+ // init tcp listener
+ listener.Init()
+ // init iptables
proxy.Init()
- go server.Start()
+ // start proxy listener
+ go listener.Start()
+ // start dns server
+ go dns.Start()
// we need watch message to update the cache of instances
for {
select {
case <-beehiveContext.Done():
klog.Warning("EdgeMesh Stop")
+ proxy.Clean()
return
default:
}
msg, err := beehiveContext.Receive(constant.ModuleNameEdgeMesh)
if err != nil {
- klog.Warningf("edgemesh receive msg error %v", err)
+ klog.Warningf("[EdgeMesh] receive msg error %v", err)
continue
}
- klog.V(4).Infof("edgemesh get message: %v", msg)
- proxy.MsgProcess(msg)
+ klog.V(4).Infof("[EdgeMesh] get message: %v", msg)
+ listener.MsgProcess(msg)
}
}
diff --git a/edgemesh/pkg/panel/fake.go b/edgemesh/pkg/plugin/panel/panel.go
index 83a1c24ea..cfebb1b35 100644
--- a/edgemesh/pkg/panel/fake.go
+++ b/edgemesh/pkg/plugin/panel/panel.go
@@ -7,30 +7,31 @@ import (
"github.com/go-chassis/go-chassis/third_party/forked/afex/hystrix-go/hystrix"
)
-type FakePanel struct {
+type EdgePanel struct {
}
-func (fp *FakePanel) GetCircuitBreaker(inv invocation.Invocation, serviceType string) (string, hystrix.CommandConfig) {
+func (ep *EdgePanel) GetCircuitBreaker(inv invocation.Invocation, serviceType string) (string, hystrix.CommandConfig) {
return "", hystrix.CommandConfig{}
}
-func (fp *FakePanel) GetLoadBalancing(inv invocation.Invocation) control.LoadBalancingConfig {
+func (ep *EdgePanel) GetLoadBalancing(inv invocation.Invocation) control.LoadBalancingConfig {
return control.LoadBalancingConfig{}
}
-func (fp *FakePanel) GetRateLimiting(inv invocation.Invocation, serviceType string) control.RateLimitingConfig {
+func (ep *EdgePanel) GetRateLimiting(inv invocation.Invocation, serviceType string) control.RateLimitingConfig {
return control.RateLimitingConfig{}
}
-func (fp *FakePanel) GetFaultInjection(inv invocation.Invocation) model.Fault {
+func (ep *EdgePanel) GetFaultInjection(inv invocation.Invocation) model.Fault {
return model.Fault{}
}
-func (fp *FakePanel) GetEgressRule() []control.EgressConfig {
+func (ep *EdgePanel) GetEgressRule() []control.EgressConfig {
return []control.EgressConfig{}
}
// TODO Remove the init method, because it will cause invalid logs to be printed when the program is running @kadisi
// init install Plugin
func init() {
- control.InstallPlugin("fake", func(options control.Options) control.Panel {
- return &FakePanel{}
+ control.InstallPlugin("edge", func(options control.Options) control.Panel {
+ return &EdgePanel{}
})
+
}
diff --git a/edgemesh/pkg/plugin/plugin.go b/edgemesh/pkg/plugin/plugin.go
new file mode 100644
index 000000000..f6e01240d
--- /dev/null
+++ b/edgemesh/pkg/plugin/plugin.go
@@ -0,0 +1,48 @@
+package plugin
+
+import (
+ "github.com/go-chassis/go-archaius"
+ "github.com/go-chassis/go-chassis/control"
+ "github.com/go-chassis/go-chassis/core/config"
+ "github.com/go-chassis/go-chassis/core/config/model"
+ "github.com/go-chassis/go-chassis/core/loadbalancer"
+ "github.com/go-chassis/go-chassis/core/registry"
+
+ meshConfig "github.com/kubeedge/kubeedge/edgemesh/pkg/config"
+ _ "github.com/kubeedge/kubeedge/edgemesh/pkg/plugin/panel"
+ meshRegistry "github.com/kubeedge/kubeedge/edgemesh/pkg/plugin/registry"
+)
+
+// Install installs go-chassis plugins
+func Install() {
+ // service discovery
+ opt := registry.Options{}
+ registry.DefaultServiceDiscoveryService = meshRegistry.NewEdgeServiceDiscovery(opt)
+ // load balance
+ loadbalancer.InstallStrategy(meshConfig.Config.LBStrategy, func() loadbalancer.Strategy {
+ switch meshConfig.Config.LBStrategy {
+ case loadbalancer.StrategyRoundRobin:
+ return &loadbalancer.RoundRobinStrategy{}
+ case loadbalancer.StrategyRandom:
+ return &loadbalancer.RandomStrategy{}
+ case loadbalancer.StrategySessionStickiness:
+ return &loadbalancer.SessionStickinessStrategy{}
+ default:
+ return &loadbalancer.RoundRobinStrategy{}
+ }
+ })
+ // control panel
+ config.GlobalDefinition = &model.GlobalCfg{
+ Panel: model.ControlPanel{
+ Infra: "edge",
+ },
+ Ssl: make(map[string]string),
+ }
+ opts := control.Options{
+ Infra: config.GlobalDefinition.Panel.Infra,
+ Address: config.GlobalDefinition.Panel.Settings["address"],
+ }
+ control.Init(opts)
+ // init archaius
+ archaius.Init()
+}
diff --git a/edgemesh/pkg/plugin/registry/registry.go b/edgemesh/pkg/plugin/registry/registry.go
new file mode 100644
index 000000000..e3135da84
--- /dev/null
+++ b/edgemesh/pkg/plugin/registry/registry.go
@@ -0,0 +1,204 @@
+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"
+
+ "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, &registry.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) (name, namespace string, port int, err error) {
+ name, namespace = common.SplitServiceKey(serviceUrl)
+ 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
+ }
+ } else {
+ klog.Errorf("[EdgeMesh] service url %s invalid", serviceUrl)
+ err = fmt.Errorf("service url %s invalid", serviceUrl)
+ }
+ return
+}
+
+// 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)
+ }
+ 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
+} \ No newline at end of file
diff --git a/edgemesh/pkg/protocol/http/http.go b/edgemesh/pkg/protocol/http/http.go
new file mode 100644
index 000000000..a38738c29
--- /dev/null
+++ b/edgemesh/pkg/protocol/http/http.go
@@ -0,0 +1,131 @@
+package http
+
+import (
+ "bufio"
+ "bytes"
+ "context"
+ "fmt"
+ "io"
+ "net"
+ "net/http"
+
+ "github.com/go-chassis/go-chassis/core/common"
+ "github.com/go-chassis/go-chassis/core/handler"
+ "github.com/go-chassis/go-chassis/core/invocation"
+ "k8s.io/klog"
+
+ "github.com/kubeedge/kubeedge/edgemesh/pkg/config"
+)
+
+type HTTP struct {
+ Conn net.Conn
+ SvcNamespace string
+ SvcName string
+ Port int
+
+ req *http.Request
+}
+
+// Process handles http protocol
+func (p *HTTP) Process() {
+ defer p.Conn.Close()
+
+ for {
+ // parse http request
+ req, err := http.ReadRequest(bufio.NewReader(p.Conn))
+ if err != nil {
+ if err == io.EOF {
+ klog.Infof("[EdgeMesh] http client disconnected.")
+ return
+ }
+ klog.Errorf("[EdgeMesh] parse http request err: %v", err)
+ return
+ }
+
+ // http: Request.RequestURI can't be set in client requests
+ // just reset it before transport
+ req.RequestURI = ""
+
+ // create invocation
+ inv := invocation.New(context.Background())
+
+ // set invocation
+ inv.MicroServiceName = req.Host
+ inv.SourceServiceID = ""
+ inv.Protocol = "rest"
+ inv.Strategy = config.Config.LBStrategy
+ inv.Args = req
+ inv.Reply = &http.Response{}
+
+ // create handlerchain
+ c, err := handler.CreateChain(common.Consumer, "http", handler.Loadbalance, handler.Transport)
+ if err != nil {
+ klog.Errorf("[EdgeMesh] create http handlerchain error: %v", err)
+ return
+ }
+
+ // start to handle
+ p.req = req
+ c.Next(inv, p.responseCallback)
+ }
+}
+
+// responseCallback implements http handlerchain callback
+func (p *HTTP) responseCallback(data *invocation.Response) error {
+ var err error
+ var resp *http.Response
+ var respBytes []byte
+
+ // as a proxy server, make sure that edgemesh always response to a request
+ // send either the response of the real backend server or 503 back
+ if data.Err != nil {
+ klog.Errorf("[EdgeMesh] error in http handlerchain: %v", data.Err)
+ err = data.Err
+ } else {
+ if data.Result == nil {
+ klog.Errorf("[EdgeMesh] empty response from http handlerchain")
+ err = fmt.Errorf("empty response from http handlerchain")
+ } else {
+ var ok bool
+ resp, ok = data.Result.(*http.Response)
+ if !ok {
+ klog.Errorf("[EdgeMesh] http handlerchain result %+v not *http.Response type", data.Result)
+ err = fmt.Errorf("result not *http.Response type")
+ } else {
+ respBytes, err = httpResponseToBytes(resp)
+ if err != nil {
+ klog.Errorf("[EdgeMesh] convert http response to bytes err: %v", err)
+ } else {
+ // send response back
+ p.Conn.Write(respBytes)
+ return nil
+ }
+ }
+ }
+ }
+ // 503
+ resp = &http.Response{
+ Status: fmt.Sprintf("%d %s", http.StatusServiceUnavailable, err),
+ StatusCode: http.StatusServiceUnavailable,
+ Proto: p.req.Proto,
+ Request: p.req,
+ Header: make(http.Header, 0),
+ }
+ respBytes, _ = httpResponseToBytes(resp)
+ // send error response back
+ p.Conn.Write(respBytes)
+ return err
+}
+
+// httpResponseToBytes transforms http.Response to bytes
+func httpResponseToBytes(resp *http.Response) ([]byte, error) {
+ buf := new(bytes.Buffer)
+ if resp == nil {
+ return nil, fmt.Errorf("http response nil")
+ }
+ err := resp.Write(buf)
+ if err != nil {
+ return nil, err
+ }
+ return buf.Bytes(), nil
+}
diff --git a/edgemesh/pkg/protocol/protocol.go b/edgemesh/pkg/protocol/protocol.go
new file mode 100644
index 000000000..8e22bc559
--- /dev/null
+++ b/edgemesh/pkg/protocol/protocol.go
@@ -0,0 +1,5 @@
+package protocol
+
+type Protocol interface {
+ Process()
+}
diff --git a/edgemesh/pkg/proxy/poll/poll.go b/edgemesh/pkg/proxy/poll/poll.go
deleted file mode 100644
index e8fea6736..000000000
--- a/edgemesh/pkg/proxy/poll/poll.go
+++ /dev/null
@@ -1,85 +0,0 @@
-package poll
-
-import (
- "syscall"
-
- "k8s.io/klog"
-)
-
-type Epoll struct {
- fd int
- callback func(fd int32)
-}
-
-const (
- eventNum = 128
- waitForever = -1
-)
-
-// CreatePoll open a epoll file description
-func CreatePoll(fn func(fd int32)) (*Epoll, error) {
- ep := &Epoll{
- callback: fn,
- }
- fd, err := syscall.EpollCreate1(0)
- if err != nil {
- return nil, err
- }
- ep.fd = fd
- return ep, err
-}
-
-// Loop for epoll
-func (ep *Epoll) Loop() {
- events := make([]syscall.EpollEvent, eventNum)
- for {
- num, err := syscall.EpollWait(ep.fd, events, waitForever)
- if num < 0 || err != nil {
- klog.Warningf("[L4 Proxy] poll wait error : %s", err)
- }
- for _, ev := range events {
- if ev.Events&syscall.EPOLLIN == syscall.EPOLLIN {
- ep.callback(ev.Fd)
- }
- }
- }
-}
-
-// DestroyEpoll implement release the epoll file description
-func (ep *Epoll) DestroyEpoll() {
- syscall.Close(ep.fd)
-}
-
-// EpollCtrlAdd: Add a socket to the epoll
-func (ep *Epoll) EpollCtrlAdd(fd int32) error {
- err := syscall.EpollCtl(ep.fd, syscall.EPOLL_CTL_ADD, int(fd), &syscall.EpollEvent{
- Fd: int32(fd),
- Events: syscall.EPOLLIN,
- })
- if err != nil {
- klog.Errorf("Add: %s , %d, %d", err, ep.fd, fd)
- }
-
- return err
-}
-
-// EpollCtrlDel: delete a socket to the epoll
-func (ep *Epoll) EpollCtrlDel(fd int) error {
- err := syscall.EpollCtl(ep.fd, syscall.EPOLL_CTL_DEL, fd, &syscall.EpollEvent{
- Fd: int32(fd),
- Events: uint32(syscall.EPOLLIN),
- })
- if err != nil {
- klog.Errorf("Del: %s , %d, %d", err, ep.fd, fd)
- }
- return err
-}
-
-// EpollCtrlMod: modify a socket to the epoll
-func (ep *Epoll) EpollCtrlMod(fd int) error {
- err := syscall.EpollCtl(ep.fd, syscall.EPOLL_CTL_MOD, fd, &syscall.EpollEvent{
- Fd: int32(fd),
- Events: syscall.EPOLLIN,
- })
- return err
-}
diff --git a/edgemesh/pkg/proxy/proxy.go b/edgemesh/pkg/proxy/proxy.go
index 50a221207..dabc128cd 100644
--- a/edgemesh/pkg/proxy/proxy.go
+++ b/edgemesh/pkg/proxy/proxy.go
@@ -1,568 +1,253 @@
package proxy
import (
- "encoding/json"
+ "bufio"
"fmt"
- "math/rand"
- "net"
- "strconv"
+ "io/ioutil"
+ "os"
"strings"
- "sync"
- "syscall"
"time"
- "github.com/kubeedge/beehive/pkg/core/model"
- "github.com/kubeedge/kubeedge/common/constants"
- "github.com/kubeedge/kubeedge/edge/pkg/metamanager/client"
- "github.com/kubeedge/kubeedge/edgemesh/pkg/proxy/poll"
- vdev "github.com/kubeedge/kubeedge/edgemesh/pkg/proxy/virtualdevice"
- v1 "k8s.io/api/core/v1"
+ "github.com/vishvananda/netlink"
"k8s.io/klog"
-)
-
-type conntrack struct {
- lconn net.Conn
- rconn net.Conn
- connNum uint32
-}
-
-type listener struct {
- ln net.Listener
- serviceName string
- fd int32
-}
+ utiliptables "k8s.io/kubernetes/pkg/util/iptables"
+ utilexec "k8s.io/utils/exec"
-type serviceTable struct {
- ip string
- ports []int32
- targetPort []int32
- lns []net.Listener
-}
+ "github.com/kubeedge/kubeedge/edgemesh/pkg/config"
+)
-type addrTable struct {
- sync.Map
+// iptables rules
+type Proxier struct {
+ iptables utiliptables.Interface
+ inboundRule string
+ outboundRule string
+ dNatRule string
}
const (
- tcpBufSize = 8192
- defaultNetworkPrefix = "9.251."
- defaultIpPoolSize = 20
- defaultTcpClientTimeout = time.Second * 2
- defaultTcpReconnectTimes = 3
- firstPort = 0
+ meshChain = "EDGE-MESH"
+ hostResolv = "/etc/resolv.conf"
)
var (
- epoll *poll.Epoll
- addrByService *addrTable
- unused []string
- serve sync.Map
- metaClient client.CoreInterface
- ipPoolSize uint16
+ proxier *Proxier
+ route netlink.Route
)
-// Init: init the proxy. create virtual device and assign ips, etc..
func Init() {
- go func() {
- unused = make([]string, 0)
- addrByService = &addrTable{}
- metaClient = client.New()
- //create virtual network device
- for {
- err := vdev.CreateDevice()
- if err == nil {
- break
- }
- klog.Warningf("[L4 Proxy] create Device is failed : %s", err)
- //there may have some exception need to be fixed on OS
- time.Sleep(2 * time.Second)
- }
- //configure vir ip
- ipPoolSize = 0
- expandIpPool()
- //open epoll
- ep, err := poll.CreatePoll(pollCallback)
- if err != nil {
- vdev.DestroyDevice()
- klog.Errorf("[L4 Proxy] epoll is open failed : %s", err)
- return
- }
- epoll = ep
- go epoll.Loop()
- klog.Infof("[L4 Proxy] proxy is running now")
- }()
-}
-
-// expandIpPool expand ip pool if virtual ip is not enough
-func expandIpPool() {
- // make sure the subnet can't be "255.255"
- if ipPoolSize > 0xfffe {
- return
+ protocol := utiliptables.ProtocolIpv4
+ exec := utilexec.New()
+ iptInterface := utiliptables.New(exec, protocol)
+ proxier = &Proxier{
+ iptables: iptInterface,
+ inboundRule: "-p tcp -d " + config.Config.SubNet + " -i " + config.Config.ListenInterface + " -j " + meshChain,
+ outboundRule: "-p tcp -d " + config.Config.SubNet + " -o " + config.Config.ListenInterface + " -j " + meshChain,
+ dNatRule: "-p tcp -j DNAT --to-destination " + config.Config.Listener.Addr().String(),
}
- for idx := ipPoolSize + 1; idx <= ipPoolSize+defaultIpPoolSize; idx++ {
- ip := defaultNetworkPrefix + getSubNet(idx)
- if err := vdev.AddIP(ip); err != nil {
- vdev.DestroyDevice()
- klog.Errorf("[L4 Proxy] Add ip is failed : %s please checkout the env", err)
- return
- }
- unused = append(unused, ip)
- }
- ipPoolSize = defaultIpPoolSize
-}
-
-// pollCallback process the connection from client
-func pollCallback(fd int32) {
- value, ok := serve.Load(fd)
- if !ok {
+ // read and clean iptables rules
+ proxier.readAndCleanRule()
+ // ensure iptables rules
+ proxier.ensureRule()
+ // add route
+ dst, err := netlink.ParseIPNet(config.Config.SubNet)
+ if err != nil {
+ klog.Errorf("[EdgeMesh] parse subnet error: %v", err)
return
}
- listen, ok := value.(listener)
- if !ok {
- return
+ gw := config.Config.ListenIP
+ route = netlink.Route{
+ Dst: dst,
+ Gw: gw,
}
- ln := listen.ln
- serviceName := listen.serviceName
-
- conn, err := ln.Accept()
+ err = netlink.RouteAdd(&route)
if err != nil {
- return
+ klog.Warningf("[EdgeMesh] add route err: %v", err)
+ }
+ // save iptables rules
+ proxier.saveRule()
+ // ensure resolv.conf
+ ensureResolvForHost()
+ // sync
+ go proxier.sync()
+}
+
+// sync periodically
+func (p *Proxier) sync() {
+ syncRuleTicker := time.NewTicker(10 * time.Second)
+ for {
+ <-syncRuleTicker.C
+ p.ensureRule()
+ ensureResolvForHost()
}
- getAndSetSocket(conn, false)
- go startTcpServer(conn, serviceName)
}
-// startTcpServer implement L4 proxy to the real server
-func startTcpServer(conn net.Conn, svcName string) {
- portString := strings.Split(conn.LocalAddr().String(), ":")
- // portString is a standard form, such as "172.17.0.1:8080"
- localPort, _ := strconv.ParseInt(portString[1], 10, 32)
- addr, err := doLoadBalance(svcName, localPort)
+// ensureRule ensures iptables rules exist
+func (p *Proxier) ensureRule() {
+ iptInterface := p.iptables
+ inboundRule := strings.Split(p.inboundRule, " ")
+ outboundRule := strings.Split(p.outboundRule, " ")
+ dNatRule := strings.Split(p.dNatRule, " ")
+ exist, err := iptInterface.EnsureChain(utiliptables.TableNAT, meshChain)
if err != nil {
- klog.Warningf("[L4 Proxy] %s call svc : %s encountered an error: %s", conn.RemoteAddr().String(), svcName, err)
- conn.Close()
- return
+ klog.Errorf("[EdgeMesh] ensure chain %s failed with err: %v", meshChain, err)
}
- var proxyClient net.Conn
- for retry := 0; retry < defaultTcpReconnectTimes; retry++ {
- proxyClient, err = net.DialTimeout("tcp", addr.String(), defaultTcpClientTimeout)
- if err == nil {
- break
- }
+ if !exist {
+ klog.Infof("[EdgeMesh] chain %s not exists", meshChain)
}
- // Error when connecting to server ,maybe timeout or any other error
+
+ exist, err = iptInterface.EnsureRule(utiliptables.Append, utiliptables.TableNAT, utiliptables.ChainPrerouting, inboundRule...)
if err != nil {
- klog.Warningf("[L4 Proxy] %s call svc : %s to %s encountered an error: %s", conn.RemoteAddr().String(), svcName, addr.String(), err)
- conn.Close()
- return
+ klog.Errorf("[EdgeMesh] ensure inbound rule %s failed with err: %v", p.inboundRule, err)
}
-
- ctk := &conntrack{
- lconn: conn,
- rconn: proxyClient,
+ if !exist {
+ klog.Infof("[EdgeMesh] inbound rule %s not exists", p.inboundRule)
}
- klog.Infof("[L4 Proxy] start a proxy server : %s,%s", svcName, addr.String())
- go func() {
- ctk.processServerProxy()
- }()
- go func() {
- ctk.processClientProxy()
- }()
-}
-// doLoadBalance implement the loadbalance function
-func doLoadBalance(svcName string, lport int64) (net.Addr, error) {
- svc := strings.Split(svcName, ".")
- namespace, name := svc[0], svc[1]
- pods, err := metaClient.Services(namespace).GetPods(name)
+ exist, err = iptInterface.EnsureRule(utiliptables.Append, utiliptables.TableNAT, utiliptables.ChainOutput, outboundRule...)
if err != nil {
- klog.Errorf("[L4 Proxy] get svc error : %s", err)
+ klog.Errorf("[EdgeMesh] ensure outbound rule %s failed with err: %v", p.outboundRule, err)
}
- // checkout the status of pods
- runPods := make([]v1.Pod, 0, len(pods))
- for i := 0; i < len(pods); i++ {
- if pods[i].Status.Phase == v1.PodRunning {
- runPods = append(runPods, pods[i])
- }
+ if !exist {
+ klog.Infof("[EdgeMesh] outbound rule %s not exists", p.outboundRule)
}
- // support random LB for the early version
- rand.Seed(time.Now().UnixNano())
- idx := rand.Uint32() % uint32(len(runPods))
-
- hostIP := runPods[idx].Status.HostIP
- st := addrByService.getAddrTable(svcName)
- index := 0
- for i, p := range st.ports {
- if p == int32(lport) {
- index = i
- break
- }
+ exist, err = iptInterface.EnsureRule(utiliptables.Append, utiliptables.TableNAT, meshChain, dNatRule...)
+ if err != nil {
+ klog.Errorf("[EdgeMesh] ensure dnat rule %s failed with err: %v", p.dNatRule, err)
}
- // kubeedge edgemesh support bridge net ,so, use hostport to access
- tp := st.targetPort[index]
- targetPort := int32(0)
- for _, value := range runPods[idx].Spec.Containers {
- for _, v := range value.Ports {
- if v.ContainerPort == tp {
- targetPort = v.HostPort
- break
- }
- }
+ if !exist {
+ klog.Infof("[EdgeMesh] dnat rule %s not exists", p.dNatRule)
}
-
- return &net.TCPAddr{
- IP: net.ParseIP(hostIP),
- Port: int(targetPort),
- }, nil
}
-// GetServiceServer returns the proxy IP by given service name
-func GetServiceServer(svcName string) string {
- st := addrByService.getAddrTable(svcName)
- if st == nil {
- klog.Warningf("[L4 Proxy] Serivce %s is not ready for Proxy.", svcName)
- return "Proxy-abnormal"
+// saveRule saves iptables rules into file
+func (p *Proxier) saveRule() {
+ file, err := os.OpenFile("/run/edgemesh-iptables", os.O_WRONLY|os.O_CREATE|os.O_TRUNC, 0666)
+ if err != nil {
+ klog.Errorf("[EdgeMesh] open file /run/edgemesh-iptables err: %v", err)
+ return
}
- return st.ip
+ // store
+ defer file.Close()
+ w := bufio.NewWriter(file)
+ fmt.Fprintln(w, p.inboundRule)
+ fmt.Fprintln(w, p.dNatRule)
+ fmt.Fprintln(w, p.outboundRule)
+ w.Flush()
}
-//getSubNet Implement uint16 convert to "uint8.uint8"
-func getSubNet(subNet uint16) string {
- arg1 := uint64(subNet & 0x00ff)
- arg2 := uint64((subNet & 0xff00) >> 8)
- return strconv.FormatUint(arg2, 10) + "." + strconv.FormatUint(arg1, 10)
-}
+// readAndCleanRule reads iptables rules from file and cleans them
+func (p *Proxier) readAndCleanRule() {
+ file, err := os.OpenFile("/run/edgemesh-iptables", os.O_RDONLY, 0444)
+ if err != nil {
+ klog.Errorf("[EdgeMesh] open file /run/edgemesh-iptables err: %v", err)
+ return
+ }
-//getAndSetSocket get file description and set socket blocking
-func getAndSetSocket(ln interface{}, nonblock bool) int {
- fd := int(-1)
- switch network := ln.(type) {
- case *net.TCPListener:
- file, err := network.File()
- if err != nil {
- klog.Infof("[L4 Proxy] get fd %s", err)
- } else {
- fd = int(file.Fd())
- }
- case *net.TCPConn:
- file, err := network.File()
- if err != nil {
- klog.Infof("[L4 Proxy] get fd %s", err)
- } else {
- fd = int(file.Fd())
+ defer file.Close()
+ scan := bufio.NewScanner(file)
+ scan.Split(bufio.ScanLines)
+ for scan.Scan() {
+ serverString := scan.Text()
+ if strings.Contains(serverString, "-o") {
+ p.iptables.DeleteRule(utiliptables.TableNAT, utiliptables.ChainOutput, strings.Split(serverString, " ")...)
+ } else if strings.Contains(serverString, "-i") {
+ p.iptables.DeleteRule(utiliptables.TableNAT, utiliptables.ChainPrerouting, strings.Split(serverString, " ")...)
}
- default:
- klog.Infof("[L4 Proxy] unknow conn")
}
+ p.iptables.FlushChain(utiliptables.TableNAT, meshChain)
+ p.iptables.DeleteChain(utiliptables.TableNAT, meshChain)
+}
- err := syscall.SetNonblock(fd, nonblock)
+// ensureResolvForHost adds edgemesh dns server to the head of /etc/resolv.conf
+func ensureResolvForHost() {
+ bs, err := ioutil.ReadFile(hostResolv)
if err != nil {
- klog.Errorf("[L4 Proxy] Set Nonblock : %s", err)
+ klog.Errorf("[EdgeMesh] read file %s err: %v", hostResolv, err)
+ return
}
- return fd
-}
-
-// addAddrTable is a thread-safe operation to add to map
-func (at *addrTable) addAddrTable(key string, value *serviceTable) {
- at.Store(key, value)
-}
-
-// addAddrTable is a thread-safe operation to del from map
-func (at *addrTable) delAddrTable(key string) {
- at.Delete(key)
-}
-
-// addAddrTable is a thread-safe operation to get from map
-func (at *addrTable) getAddrTable(key string) *serviceTable {
- value, ok := at.Load(key)
- if !ok {
- return nil
- }
- st, ok := value.(*serviceTable)
- if !ok {
- return nil
+ resolv := strings.Split(string(bs), "\n")
+ if resolv == nil {
+ nameserver := "nameserver " + config.Config.ListenIP.String()
+ ioutil.WriteFile(hostResolv, []byte(nameserver), 0600)
+ return
}
- return st
-}
-// filterResourceType implement filter. Proxy cares "Service" and "ServiceList" type
-func filterResourceType(msg model.Message) []v1.Service {
- svcs := make([]v1.Service, 0)
- var content []byte
- var err error
- switch msg.Content.(type) {
- case []byte:
- content = msg.GetContent().([]byte)
- default:
- content, err = json.Marshal(msg.GetContent())
- if err != nil {
- klog.Errorf("marshal message to edgemesh failed, err: %v", err)
- }
- }
- switch getReourceType(msg.GetResource()) {
- case constants.ResourceTypeService:
- s, err := handleServiceMessage(content)
- if err != nil {
+ configured := false
+ dnsIdx := 0
+ startIdx := 0
+ for idx, item := range resolv {
+ if strings.Contains(item, config.Config.ListenIP.String()) {
+ configured = true
+ dnsIdx = idx
break
}
- svcs = append(svcs, *s)
- case constants.ResourceTypeServiceList:
- ss, err := handleServiceMessageList(content)
- if err != nil {
+ }
+ for idx, item := range resolv {
+ if strings.Contains(item, "nameserver") {
+ startIdx = idx
break
}
- svcs = append(svcs, ss...)
- default:
- klog.Infof("[L4 Proxy] process other resource: %s", msg.Router.Resource)
}
-
- return svcs
-}
-
-// MsgProcess process from metaManager and start a proxy server
-func MsgProcess(msg model.Message) {
- svcs := filterResourceType(msg)
- if len(svcs) == 0 {
+ if configured {
+ if dnsIdx != startIdx && dnsIdx > startIdx {
+ nameserver := sortNameserver(resolv, dnsIdx, startIdx)
+ ioutil.WriteFile(hostResolv, []byte(nameserver), 0600)
+ }
return
}
- klog.Infof("[L4 Proxy] proxy process svcs : %d resource: %s\n", len(svcs), msg.Router.Resource)
- for _, svc := range svcs {
- svcName := svc.Namespace + "." + svc.Name
- if !IsL4Proxy(&svc) {
- // when server protocol update to http
- delServer(svcName)
- continue
- }
- klog.Infof("[L4 Proxy] proxy process svc : %s,%s", msg.GetOperation(), svcName)
- port := make([]int32, 0)
- targetPort := make([]int32, 0)
- for _, p := range svc.Spec.Ports {
- // this version will support TCP only
- if p.Protocol == "TCP" {
- port = append(port, p.Port)
- // this version will not support string type
- targetPort = append(targetPort, p.TargetPort.IntVal)
- }
- }
- if len(port) == 0 || len(targetPort) == 0 {
+ nameserver := ""
+ for idx := 0; idx < len(resolv); {
+ if idx == startIdx {
+ startIdx = -1
+ nameserver = nameserver + "nameserver " + config.Config.ListenIP.String() + "\n"
continue
}
- switch msg.GetOperation() {
- case "insert":
- addServer(svcName, port)
- case "delete":
- delServer(svcName)
- case "update":
- updateServer(svcName, port)
- default:
- klog.Infof("[L4 proxy] Unknown operation")
- }
- st := addrByService.getAddrTable(svcName)
- if st != nil {
- st.targetPort = targetPort
- }
- }
-}
-
-// addServer : add the proxy server
-func addServer(svcName string, ports []int32) {
- var ip string
- st := addrByService.getAddrTable(svcName)
- if st != nil {
- if len(ports) == 0 && len(st.ports) == 0 {
- unused = append(unused, st.ip)
- addrByService.delAddrTable(svcName)
- return
- }
- ip = st.ip
- } else {
- if len(ports) == 0 {
- return
- }
- if len(unused) == 0 {
- expandIpPool()
- }
- ip = unused[0]
- unused = unused[1:]
+ nameserver = nameserver + resolv[idx] + "\n"
+ idx++
}
- lns := make([]net.Listener, 0)
- for _, port := range ports {
- addr := ip + ":" + strconv.FormatUint(uint64(port), 10)
- klog.Infof("[L4 Proxy] Start listen %s,%d for proxy", addr, port)
- ln, err := net.Listen("tcp", addr)
- if err != nil {
- klog.Errorf("[L4 Proxy] %s", err)
- continue
- }
- lns = append(lns, ln)
- server := listener{
- ln: ln,
- serviceName: svcName,
- fd: int32(getAndSetSocket(ln, true)),
- }
- serve.Store(server.fd, server)
- epoll.EpollCtrlAdd(server.fd)
- }
- if st != nil {
- st.lns = append(st.lns, lns...)
- st.ports = append(st.ports, ports...)
- } else {
- st = &serviceTable{
- ip: ip,
- lns: lns,
- ports: ports,
- }
- addrByService.addAddrTable(svcName, st)
- }
+ ioutil.WriteFile(hostResolv, []byte(nameserver), 0600)
}
-// delServer implement delete the proxy server
-func delServer(svcName string) {
- st := addrByService.getAddrTable(svcName)
-
- if st == nil {
- return
+func sortNameserver(resolv []string, dnsIdx, startIdx int) string {
+ nameserver := ""
+ idx := 0
+ for ; idx < startIdx; idx++ {
+ nameserver = nameserver + resolv[idx] + "\n"
}
- unused = append(unused, st.ip)
- for _, ln := range st.lns {
- fd := getAndSetSocket(ln, true)
- epoll.EpollCtrlDel(fd)
- ln.Close()
- }
- addrByService.delAddrTable(svcName)
-}
+ nameserver = nameserver + resolv[dnsIdx] + "\n"
-// updateServer implement update the proxy server
-func updateServer(svcName string, ports []int32) error {
- st := addrByService.getAddrTable(svcName)
- if st == nil { // if not exist
- addServer(svcName, ports)
- } else {
- oldports := make([]int32, len(st.ports))
- copy(oldports, st.ports)
- for idx, oldport := range oldports {
- update := true
- for k, newport := range ports {
- if oldport == newport {
- ports = append(ports[:k], ports[k+1:]...)
- update = false
- break
- }
- }
- if update {
- fd := getAndSetSocket(st.lns[idx], true)
- epoll.EpollCtrlDel(fd)
- st.lns[idx].Close()
- st.lns = append(st.lns[:idx], st.lns[idx+1:]...)
- st.ports = append(st.ports[:idx], st.ports[idx+1:]...)
- }
+ for idx = startIdx; idx < len(resolv); idx++ {
+ if idx == dnsIdx {
+ continue
}
- addServer(svcName, ports)
+ nameserver = nameserver + resolv[idx] + "\n"
}
- return nil
-}
-//handleMessageFromMetaManager convert []byte to k8s Service struct
-func handleServiceMessage(content []byte) (*v1.Service, error) {
- var s v1.Service
- err := json.Unmarshal(content, &s)
- if err != nil {
- return nil, fmt.Errorf("[L4 Proxy] unmarshal message to Service failed, err: %v", err)
- }
- return &s, nil
+ return nameserver
}
-//handleMessageFromMetaManager convert []byte to k8s Service struct
-func handleServiceMessageList(content []byte) ([]v1.Service, error) {
- var s []v1.Service
- err := json.Unmarshal(content, &s)
+func Clean() {
+ proxier.readAndCleanRule()
+ netlink.RouteDel(&route)
+ bs, err := ioutil.ReadFile(hostResolv)
if err != nil {
- return nil, fmt.Errorf("[L4 Proxy] unmarshal message to Service failed, err: %v", err)
+ klog.Warningf("[EdgeMesh] read file %s err: %v", hostResolv, err)
}
- return s, nil
-}
-
-// getReourceType returns the reourceType as a string
-func getReourceType(reource string) string {
- str := strings.Split(reource, "/")
- if len(str) == 3 {
- return str[1]
- } else if len(str) == 5 {
- return str[3]
- } else {
- return reource
- }
-}
-// processServerProxy process up link traffic
-func (c *conntrack) processClientProxy() {
- buf := make([]byte, tcpBufSize)
- for {
- n, err := c.lconn.Read(buf)
- //service caller closess the connection
- if n == 0 {
- c.lconn.Close()
- c.rconn.Close()
- break
- }
- if err != nil {
- c.lconn.Close()
- c.rconn.Close()
- break
- }
- _, rerr := c.rconn.Write(buf[:n])
- if rerr != nil {
- c.lconn.Close()
- c.rconn.Close()
- break
- }
+ resolv := strings.Split(string(bs), "\n")
+ if resolv == nil {
+ return
}
-}
-
-// processServerProxy process down link traffic
-func (c *conntrack) processServerProxy() {
- buf := make([]byte, tcpBufSize)
- for {
- n, err := c.rconn.Read(buf)
- if n == 0 {
- c.rconn.Close()
- c.lconn.Close()
- break
- }
- if err != nil {
- c.rconn.Close()
- c.lconn.Close()
- break
- }
- _, rerr := c.lconn.Write(buf[:n])
- if rerr != nil {
- c.rconn.Close()
- c.lconn.Close()
- break
+ nameserver := ""
+ for _, item := range resolv {
+ if strings.Contains(item, config.Config.ListenIP.String()) {
+ continue
}
+ nameserver = nameserver + item + "\n"
}
-}
-
-//isL4Proxy Determine whether to use L4 proxy
-func IsL4Proxy(svc *v1.Service) bool {
- if len(svc.Spec.Ports) == 0 {
- return false
- }
- // In the defination of k8s-Service, we can use Service.Spec.Ports.Name to
- // indicate whether the Service enables L4 proxy mode. According to our
- // current L7 mode only support the http protocol. Other 7-layer protocols
- // are automatically degraded to tcp until supported
- port := svc.Spec.Ports[firstPort]
- switch port.Name {
- case "websocket", "grpc", "https", "tcp":
- return true
- case "http", "udp":
- return false
- default:
- return true
- }
+ ioutil.WriteFile(hostResolv, []byte(nameserver), 0600)
}
diff --git a/edgemesh/pkg/proxy/virtualdevice/virtualdevice.go b/edgemesh/pkg/proxy/virtualdevice/virtualdevice.go
deleted file mode 100644
index c9bfceb20..000000000
--- a/edgemesh/pkg/proxy/virtualdevice/virtualdevice.go
+++ /dev/null
@@ -1,72 +0,0 @@
-package virtualdevice
-
-import (
- "fmt"
- "net"
-
- "github.com/vishvananda/netlink"
- "golang.org/x/sys/unix"
-)
-
-const (
- DeviceNameDefault = "edge0"
-)
-
-var (
- nh *netlink.Handle
-)
-
-func CreateDevice() error {
- //if device is exist,delete and create a new one
- //make sure the latest configuration
- DestroyDevice()
-
- edge0 := &netlink.Dummy{
- LinkAttrs: netlink.LinkAttrs{
- Name: DeviceNameDefault,
- },
- }
- return nh.LinkAdd(edge0)
-}
-
-// AddIP Bind ip to the default device
-func AddIP(address string) error {
- dev, err := nh.LinkByName(DeviceNameDefault)
- if err != nil {
- return fmt.Errorf("[L4 Proxy] Device %s is not exist!!", DeviceNameDefault)
- }
-
- addr := net.ParseIP(address)
- if addr == nil {
- return fmt.Errorf("[L4 Proxy] AddIP address is invalid ")
- }
-
- if err := nh.AddrAdd(dev, &netlink.Addr{IPNet: netlink.NewIPNet(addr)}); err != nil {
- if err == unix.EEXIST {
- return nil
- }
- return err
- }
-
- return nil
-}
-
-// DestroyDevice implement release the device
-func DestroyDevice() error {
- _, err := net.InterfaceByName(DeviceNameDefault)
- if err != nil {
- return err
- }
-
- link, err := nh.LinkByName(DeviceNameDefault)
- edge0, ok := link.(*netlink.Dummy)
- if !ok {
- return fmt.Errorf("")
- }
-
- return nh.LinkDel(edge0)
-}
-
-func init() {
- nh = &netlink.Handle{}
-}
diff --git a/edgemesh/pkg/registry/registry.go b/edgemesh/pkg/registry/registry.go
deleted file mode 100644
index 335482946..000000000
--- a/edgemesh/pkg/registry/registry.go
+++ /dev/null
@@ -1,95 +0,0 @@
-package registry
-
-import (
- "fmt"
- "strconv"
-
- "github.com/go-chassis/go-chassis/core/registry"
- utiltags "github.com/go-chassis/go-chassis/pkg/util/tags"
- v1 "k8s.io/api/core/v1"
- "k8s.io/klog"
-
- "github.com/kubeedge/kubeedge/edge/pkg/metamanager/client"
- "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, NewServiceDiscovery) }
-
-// ServiceDiscovery to represent the object of service center to call the APIs of service center
-type ServiceDiscovery struct {
- metaClient client.CoreInterface
- Name string
-}
-
-func toProtocolMap(address v1.EndpointAddress, ports []v1.EndpointPort) map[string]string {
- ret := map[string]string{}
- for _, port := range ports {
- if _, ok := ret[port.Name]; !ok {
- ret[port.Name] = address.IP + ":" + strconv.Itoa(int(port.Port))
- continue
- }
- }
- return ret
-}
-
-func NewServiceDiscovery(options registry.Options) registry.ServiceDiscovery {
- return &ServiceDiscovery{
- metaClient: client.New(),
- Name: EdgeRegistry,
- }
-}
-
-// GetAllMicroServices Get all MicroService information.
-func (r *ServiceDiscovery) GetAllMicroServices() ([]*registry.MicroService, error) {
- return nil, nil
-}
-
-// FindMicroServiceInstances find micro-service instances (subnets)
-func (r *ServiceDiscovery) FindMicroServiceInstances(consumerID, microServiceName string, tags utiltags.Tags) ([]*registry.MicroServiceInstance, error) {
- name, namespace := common.SplitServiceKey(microServiceName)
-
- pods, err := r.metaClient.Services(namespace).GetPods(name)
- if err != nil {
- klog.Errorf("get service pod list failed, error: %v", err)
- return nil, err
- }
- var microServiceInstance []*registry.MicroServiceInstance
- for _, p := range pods {
- microServiceInstance = append(microServiceInstance, &registry.MicroServiceInstance{
- InstanceID: "",
- ServiceID: name + "." + namespace,
- HostName: "",
- EndpointsMap: map[string]string{"rest": fmt.Sprintf("%s:%d", p.Status.HostIP, p.Spec.Containers[0].Ports[0].HostPort)},
- })
- }
-
- return microServiceInstance, nil
-}
-
-// GetMicroServiceID get microServiceID
-func (r *ServiceDiscovery) GetMicroServiceID(appID, microServiceName, version, env string) (string, error) {
- return "", nil
-}
-
-// GetMicroServiceInstances return instances
-func (r *ServiceDiscovery) GetMicroServiceInstances(consumerID, providerID string) ([]*registry.MicroServiceInstance, error) {
- return nil, nil
-}
-
-// GetMicroService return service
-func (r *ServiceDiscovery) GetMicroService(microServiceID string) (*registry.MicroService, error) {
- return nil, nil
-}
-
-// AutoSync updating the cache manager
-func (r *ServiceDiscovery) AutoSync() {}
-
-// Close close all websocket connection
-func (r *ServiceDiscovery) Close() error { return nil }
diff --git a/edgemesh/pkg/resolver/resolver.go b/edgemesh/pkg/resolver/resolver.go
deleted file mode 100644
index f35da0a3a..000000000
--- a/edgemesh/pkg/resolver/resolver.go
+++ /dev/null
@@ -1,74 +0,0 @@
-package resolver
-
-import (
- "bufio"
- "bytes"
- "context"
- "net/http"
- "strings"
-
- "github.com/go-chassis/go-chassis/core/invocation"
- "k8s.io/klog"
-
- "github.com/kubeedge/kubeedge/edgemesh/pkg/config"
-)
-
-type Resolver interface {
- Resolve(chan []byte, chan interface{}, func(string, invocation.Invocation)) (invocation.Invocation, bool)
-}
-
-type MyResolver struct {
- Name string
-}
-
-func httpMethods() (methods []string) {
- methods = []string{"GET", "HEAD", "POST", "OPTIONS", "PUT", "DELETE", "TRACE", "CONNECT"}
- return
-}
-
-func isHTTPRequest(s string) bool {
- methods := httpMethods()
- for _, method := range methods {
- if strings.HasPrefix(s, method) {
- return true
- }
- }
- return false
-}
-
-func (resolver *MyResolver) Resolve(data chan []byte, stop chan interface{}, invCallback func(string, invocation.Invocation)) (invocation.Invocation, bool) {
- content := ""
- protocol := ""
- for {
- select {
- case d := <-data:
- strData := string(d[:])
- if protocol == "" {
- if isHTTPRequest(strData) {
- protocol = "http"
- } else {
- return invocation.Invocation{}, false
- }
- }
- content += strData
- req, err := http.ReadRequest(bufio.NewReader(bytes.NewReader([]byte(content))))
- if err == nil {
- content = ""
- req.RequestURI = ""
- i := invocation.New(context.Background())
- i.MicroServiceName = req.Host
- i.SourceServiceID = ""
- i.Protocol = "rest"
- i.Args = req
- i.Strategy = config.Config.LBStrategy
- i.Reply = &http.Response{}
- invCallback("http", *i)
- }
- case <-stop:
- i := invocation.Invocation{MicroServiceName: resolver.Name, Args: content}
- invCallback(protocol, i)
- return i, true
- }
- klog.Infof("content: %s\n", content)
- }
-}
diff --git a/edgemesh/pkg/resolver/resolver_chain.go b/edgemesh/pkg/resolver/resolver_chain.go
deleted file mode 100644
index 79e4d98ea..000000000
--- a/edgemesh/pkg/resolver/resolver_chain.go
+++ /dev/null
@@ -1,28 +0,0 @@
-package resolver
-
-import (
- "container/list"
-
- "github.com/go-chassis/go-chassis/core/invocation"
-)
-
-var ResolverChain *list.List
-
-func init() {
- ResolverChain = list.New()
-}
-
-// Resolve will loop the resolverchain to resolve the request
-func Resolve(request chan []byte, stop chan interface{}, invCallback func(string, invocation.Invocation)) (invocation.Invocation, bool) {
- for resolver := ResolverChain.Front(); resolver != nil; resolver = resolver.Next() {
- inv, isFired := resolver.Value.(Resolver).Resolve(request, stop, invCallback)
- if isFired {
- return inv, true
- }
- }
- return invocation.Invocation{}, false
-}
-
-func RegisterResolver(resolver Resolver) {
- ResolverChain.PushBack(resolver)
-}
diff --git a/edgemesh/pkg/resolver/resolver_chain_test.go b/edgemesh/pkg/resolver/resolver_chain_test.go
deleted file mode 100644
index 54a23a975..000000000
--- a/edgemesh/pkg/resolver/resolver_chain_test.go
+++ /dev/null
@@ -1,91 +0,0 @@
-package resolver_test
-
-import (
- "strings"
- "testing"
- "time"
-
- "github.com/go-chassis/go-chassis/core/invocation"
- "k8s.io/klog"
-
- "github.com/kubeedge/kubeedge/edgemesh/pkg/resolver"
-)
-
-type TestResolver struct {
- Name string
-}
-
-func (resolver *TestResolver) Resolve(data chan []byte, stop chan interface{}, invCallback func(string, invocation.Invocation)) (invocation.Invocation, bool) {
- content := ""
- protocol := ""
- for {
- select {
- case d := <-data:
- strData := string(d[:])
- if protocol == "" {
- //Only address HTTP
- if strings.HasPrefix(strData, resolver.Name) {
- protocol = resolver.Name
- content += strData
- } else {
- return invocation.Invocation{}, false
- }
- } else {
- content += strData
- }
- case <-stop:
- i := invocation.Invocation{MicroServiceName: resolver.Name, Args: content}
- invCallback(protocol, i)
- return i, true
- }
- klog.Infof("content: %s\n", content)
- }
-}
-
-func TestResolve(t *testing.T) {
- //Register resolver
- r1 := &TestResolver{"http"}
- r2 := &TestResolver{"grpc"}
- resolver.RegisterResolver(r1)
- resolver.RegisterResolver(r2)
- invCallback := func(protocol string, inv invocation.Invocation) {
- klog.Infof("protocol in invCallback:%v", protocol)
- klog.Infof("content in invCallback: %v\n", inv.Args)
- }
- d := make(chan []byte, 1024)
- s := make(chan interface{}, 1)
-
- //Do resolver
- go func() {
- i, f := resolver.Resolve(d, s, invCallback)
- if !f {
- t.Error("resolver chain resolve error to no able to fire an existing resolver")
- } else {
- if i.MicroServiceName != "http" {
- t.Error("resolver chain resolve error when construct invocation")
- }
- }
- }()
-
- d <- []byte("http://support.huaweicloud.com/")
- time.Sleep(10 * time.Millisecond)
- d <- []byte("usermanual-ief/ief_01_0001.html")
- time.Sleep(10 * time.Millisecond)
- close(s)
-
- d = make(chan []byte, 1024)
- s = make(chan interface{}, 1)
- go func() {
- _, f := resolver.Resolve(d, s, invCallback)
- if f {
- t.Error("resolver chain resolve error to fired a no-existing resolver")
- }
- }()
- d <- []byte("quic://support.huaweicloud.com/")
- time.Sleep(10 * time.Millisecond)
- d <- []byte("usermanual-ief/ief_01_0001.html")
- time.Sleep(10 * time.Millisecond)
- close(s)
-
- time.Sleep(3 * time.Second)
-}
diff --git a/edgemesh/pkg/server/dns.go b/edgemesh/pkg/server/dns.go
deleted file mode 100644
index 7328be424..000000000
--- a/edgemesh/pkg/server/dns.go
+++ /dev/null
@@ -1,427 +0,0 @@
-package server
-
-import (
- "bufio"
- "encoding/binary"
- "errors"
- "fmt"
- "net"
- "os"
- "strings"
- "time"
- "unsafe"
-
- "k8s.io/klog"
-
- "github.com/kubeedge/kubeedge/edge/pkg/metamanager/client"
- "github.com/kubeedge/kubeedge/edgemesh/pkg/common"
- "github.com/kubeedge/kubeedge/edgemesh/pkg/proxy"
-)
-
-var (
- dnsQr = uint16(0x8000)
- oneByteSize = uint16(1)
- twoByteSize = uint16(2)
- ttl = uint32(64)
- defaultFakeIP = []byte{5, 5, 5, 5}
-)
-
-const (
- aRecord = 1
- bufSize = 1024
- notImplem = uint16(0x0004)
- serverFAilure = uint16(0x0002)
-)
-
-type dnsHeader struct {
- transactionID uint16
- flags uint16
- queNum uint16
- ansNum uint16
- authNum uint16
- additNum uint16
-}
-
-type dnsQuestion struct {
- from *net.UDPAddr
- head *dnsHeader
- name []byte
- queByte []byte
- qType uint16
- qClasss uint16
- queNum uint16
-}
-
-type dnsAnswer struct {
- name []byte
- qType uint16
- qClass uint16
- ttl uint32
- dataLen uint16
- addr []byte
-}
-
-//define the dns Question list type
-type dnsQs []dnsQuestion
-
-//metaClient is a query client
-var metaClient client.CoreInterface
-
-//dnsConn save DNS server
-var dnsConn *net.UDPConn
-
-//DnsStart is a External interface
-func DnsStart() {
- startDnsServer()
-}
-
-// startDnsServer start the DNS Server
-func startDnsServer() {
- metaClient = client.New()
- //get DNS server name
- lip, err := getIP()
- if err != nil {
- klog.Errorf("Dns server Start error : %s", err)
- return
- }
-
- laddr := &net.UDPAddr{
- IP: lip,
- Port: 53,
- }
- udpConn, err := net.ListenUDP("udp", laddr)
- if err != nil {
- klog.Errorf("Dns server Start error : %s", err)
- return
- }
- defer udpConn.Close()
- dnsConn = udpConn
- for {
- req := make([]byte, bufSize)
- n, from, err := dnsConn.ReadFromUDP(req)
- if err != nil || n <= 0 {
- klog.Infof("DNS server get an IO error : %s", err)
- continue
- }
-
- que, err := parseDnsQuery(req[:n])
- if err != nil {
- continue
- }
- for i, _ := range que {
- que[i].from = from
- }
- rsp := make([]byte, 0)
- rsp, err = recordHandler(que, req[0:n])
- if err != nil {
- klog.Infof("DNS server get an resolve abnormal : %s", err)
- continue
- }
- dnsConn.WriteTo(rsp, from)
- }
-}
-
-//recordHandler returns the Answer for the dns question
-func recordHandler(que []dnsQuestion, req []byte) (rsp []byte, err error) {
- var exist bool
- var ip string
- for _, q := range que {
- domainName := string(q.name)
- exist, ip = lookupFromMetaManager(domainName)
- if !exist {
- //if this service don't belongs to this cluster
- go getfromRealDNS(req, q.from)
- return rsp, fmt.Errorf("get from real DNS")
- }
- }
- fakeIP := defaultFakeIP
- if ip != "" {
- fakeIP = net.ParseIP(ip).To4()
- }
-
- pre := modifyRspPrefix(que, fakeIP)
- rsp = append(rsp, pre...)
- for _, q := range que {
- // head of each que is the same
- if que[0].head.ansNum == 0 {
- continue
- }
- //create a deceptive rep
- dnsAns := &dnsAnswer{
- name: q.name,
- qType: q.qType,
- qClass: q.qClasss,
- ttl: ttl,
- dataLen: uint16(len(fakeIP)),
- addr: fakeIP,
- }
- ans := dnsAns.getAnswer()
- rsp = append(rsp, ans...)
- }
-
- return rsp, nil
-}
-
-//parseDnsQuery returns question of the dns request
-func parseDnsQuery(req []byte) (que []dnsQuestion, err error) {
- head := &dnsHeader{}
- head.getHeader(req)
- if !head.isAQurey() {
- return nil, errors.New("Ignore")
- }
-
- question := make(dnsQs, head.queNum)
-
- offset := uint16(unsafe.Sizeof(dnsHeader{}))
- question.getQuestion(req, offset, head)
-
- que = question
- err = nil
- return
-}
-
-//isAQuery judge if the dns pkg is a Qurey process
-func (h *dnsHeader) isAQurey() bool {
- return h.flags&dnsQr != dnsQr
-}
-
-//getHeader get dns pkg head
-func (h *dnsHeader) getHeader(req []byte) {
- h.transactionID = binary.BigEndian.Uint16(req[0:2])
- h.flags = binary.BigEndian.Uint16(req[2:4])
- h.queNum = binary.BigEndian.Uint16(req[4:6])
- h.ansNum = binary.BigEndian.Uint16(req[6:8])
- h.authNum = binary.BigEndian.Uint16(req[8:10])
- h.additNum = binary.BigEndian.Uint16(req[10:12])
-}
-
-//getQuestion get dns questions
-func (q dnsQs) getQuestion(req []byte, offset uint16, head *dnsHeader) {
- ost := offset
- qNum := uint16(len(q))
-
- for i := uint16(0); i < qNum; i++ {
- tmp := ost
- ost = q[i].getQName(req, ost)
- q[i].qType = binary.BigEndian.Uint16(req[ost : ost+twoByteSize])
- ost += twoByteSize
- q[i].qClasss = binary.BigEndian.Uint16(req[ost : ost+twoByteSize])
- ost += twoByteSize
- q[i].head = head
- q[i].queByte = req[tmp:ost]
- }
-}
-
-//getAnswer Generate Answer for the dns question
-func (d *dnsAnswer) getAnswer() (ans []byte) {
- ans = make([]byte, 0)
-
- if d.qType == aRecord {
- ans = append(ans, 0xc0)
- ans = append(ans, 0x0c)
-
- tmp16 := make([]byte, 2)
- tmp32 := make([]byte, 4)
-
- binary.BigEndian.PutUint16(tmp16, d.qType)
- ans = append(ans, tmp16...)
- binary.BigEndian.PutUint16(tmp16, d.qClass)
- ans = append(ans, tmp16...)
- binary.BigEndian.PutUint32(tmp32, d.ttl)
- ans = append(ans, tmp32...)
- binary.BigEndian.PutUint16(tmp16, d.dataLen)
- ans = append(ans, tmp16...)
- ans = append(ans, d.addr...)
- }
-
- return ans
-}
-
-// getQName get dns question domain name
-func (q *dnsQuestion) getQName(req []byte, offset uint16) uint16 {
- ost := offset
-
- for {
- qbyte := uint16(req[ost])
-
- if qbyte == 0x00 {
- q.name = q.name[:uint16(len(q.name))-oneByteSize]
- return ost + oneByteSize
- }
- ost += oneByteSize
- q.name = append(q.name, req[ost:ost+qbyte]...)
- q.name = append(q.name, 0x2e)
- ost += qbyte
- }
-}
-
-// lookupFromMetaManager implement confirm the service exists
-func lookupFromMetaManager(serviceUrl string) (exist bool, ip string) {
- name, namespace := common.SplitServiceKey(serviceUrl)
- s, _ := metaClient.Services(namespace).Get(name)
- if s != nil {
- ip := ""
- //Determine whether to use L4 proxy
- if proxy.IsL4Proxy(s) {
- svcName := namespace + "." + name
- ip = proxy.GetServiceServer(svcName)
- }
- klog.Infof("Service %s is found in this cluster. namespace : %s, name: %s", serviceUrl, namespace, name)
- return true, ip
- }
- klog.Infof("Service %s is not found in this cluster", serviceUrl)
- return false, ""
-}
-
-// getfromRealDNS returns the dns response from the real DNS server
-func getfromRealDNS(req []byte, from *net.UDPAddr) {
- rsp := make([]byte, 0)
- ips, err := parseNameServer()
- if err != nil {
- return
- }
-
- laddr := &net.UDPAddr{
- IP: net.IPv4zero,
- Port: 0,
- }
-
- for _, ip := range ips { // get from real
- raddr := &net.UDPAddr{
- IP: ip,
- Port: 53,
- }
- conn, err := net.DialUDP("udp", laddr, raddr)
- defer conn.Close()
- if err != nil {
- continue
- }
- _, err = conn.Write(req)
- if err != nil {
- continue
- }
- if err = conn.SetReadDeadline(time.Now().Add(time.Minute)); err != nil {
- continue
- }
- var n int
- buf := make([]byte, bufSize)
- n, err = conn.Read(buf)
- if err != nil {
- continue
- }
-
- if n > 0 {
- rsp = append(rsp, buf[:n]...)
- dnsConn.WriteToUDP(rsp, from)
- break
- }
- }
-}
-
-// parseNameServer gets the nameserver from the resolv.conf
-func parseNameServer() ([]net.IP, error) {
- file, err := os.Open("/etc/resolv.conf")
- if err != nil {
- return nil, fmt.Errorf("error opening /etc/resolv.conf : %s", err)
- }
- defer file.Close()
-
- scan := bufio.NewScanner(file)
- scan.Split(bufio.ScanLines)
-
- ip := make([]net.IP, 0)
-
- for scan.Scan() { //get name server
- serverString := scan.Text()
- fmt.Println(serverString)
- if strings.Contains(serverString, "nameserver") {
- tmpString := strings.Replace(serverString, "nameserver", "", 1)
- nameserver := strings.TrimSpace(tmpString)
- sip := net.ParseIP(nameserver)
- if sip != nil {
- ip = append(ip, sip)
- }
- }
- }
- if len(ip) == 0 {
- return nil, fmt.Errorf("there is no nameserver in /etc/resolv.conf")
- }
- return ip, nil
-}
-
-// modifyRspPrefix use req' head generate a rsp head
-func modifyRspPrefix(que []dnsQuestion, fakeip []byte) (pre []byte) {
- ansNum := len(que)
- if ansNum == 0 {
- return
- }
- //use head in que. All the same
- rspHead := que[0].head
- rspHead.converQueryRsp(true)
-
- serverError := false
- if fakeip == nil || len(fakeip) != 4 {
- serverError = true
- }
- if que[0].qType == aRecord && (!serverError) {
- rspHead.setAnswerNum(uint16(ansNum))
- } else {
- rspHead.setAnswerNum(0)
- }
-
- rspHead.setRspRcode(que, serverError)
- pre = rspHead.getByteFromDnsHeader()
-
- for _, q := range que {
- pre = append(pre, q.queByte...)
- }
-
- return
-}
-
-// converQueryRsp conversion the dns head to a response for one query
-func (h *dnsHeader) converQueryRsp(isRsp bool) {
- if isRsp {
- h.flags |= dnsQr
- } else {
- h.flags |= dnsQr
- }
-}
-
-// set the Answer num for dns head
-func (h *dnsHeader) setAnswerNum(num uint16) {
- h.ansNum = num
-}
-
-// set the dns response return code
-func (h *dnsHeader) setRspRcode(que dnsQs, serverError bool) {
- for _, q := range que {
- if q.qType != aRecord {
- h.flags &= (^notImplem)
- h.flags |= notImplem
- } else if serverError {
- h.flags &= (^serverFAilure)
- h.flags |= serverFAilure
- }
- }
-}
-
-//getByteFromDnsHeader implement from dnsHeader struct to []byte
-func (h *dnsHeader) getByteFromDnsHeader() (rspHead []byte) {
- rspHead = make([]byte, unsafe.Sizeof(*h))
-
- idxTran := unsafe.Sizeof(h.transactionID)
- idxflag := unsafe.Sizeof(h.flags) + idxTran
- idxque := unsafe.Sizeof(h.ansNum) + idxflag
- idxans := unsafe.Sizeof(h.ansNum) + idxque
- idxaut := unsafe.Sizeof(h.authNum) + idxans
- idxadd := unsafe.Sizeof(h.additNum) + idxaut
-
- binary.BigEndian.PutUint16(rspHead[:idxTran], h.transactionID)
- binary.BigEndian.PutUint16(rspHead[idxTran:idxflag], h.flags)
- binary.BigEndian.PutUint16(rspHead[idxflag:idxque], h.queNum)
- binary.BigEndian.PutUint16(rspHead[idxque:idxans], h.ansNum)
- binary.BigEndian.PutUint16(rspHead[idxans:idxaut], h.authNum)
- binary.BigEndian.PutUint16(rspHead[idxaut:idxadd], h.additNum)
- return
-}
diff --git a/edgemesh/pkg/server/ip.go b/edgemesh/pkg/server/ip.go
deleted file mode 100644
index 9d2716af3..000000000
--- a/edgemesh/pkg/server/ip.go
+++ /dev/null
@@ -1,28 +0,0 @@
-package server
-
-import (
- "net"
- "time"
-
- "k8s.io/klog"
-)
-
-const inter = "docker0"
-
-// getIP returns the specific interface ip of version 4
-func getIP() (net.IP, error) {
- for {
- ifaces, err := net.InterfaceByName(inter)
- if err != nil {
- return nil, err
- }
- addrs, _ := ifaces.Addrs()
- for _, addr := range addrs {
- if ip, inet, _ := net.ParseCIDR(addr.String()); len(inet.Mask) == 4 {
- return ip, nil
- }
- }
- klog.Warningf("the interface %s have not config ip of version 4", inter)
- time.Sleep(time.Second * 3)
- }
-}
diff --git a/edgemesh/pkg/server/server.go b/edgemesh/pkg/server/server.go
deleted file mode 100644
index 5a93166af..000000000
--- a/edgemesh/pkg/server/server.go
+++ /dev/null
@@ -1,49 +0,0 @@
-package server
-
-import (
- "github.com/go-chassis/go-chassis/control"
- "github.com/go-chassis/go-chassis/core/config"
- "github.com/go-chassis/go-chassis/core/config/model"
- "github.com/go-chassis/go-chassis/core/loadbalancer"
- "github.com/go-chassis/go-chassis/core/registry"
-
- meshconfig "github.com/kubeedge/kubeedge/edgemesh/pkg/config"
- _ "github.com/kubeedge/kubeedge/edgemesh/pkg/panel"
- edgeregistry "github.com/kubeedge/kubeedge/edgemesh/pkg/registry"
- "github.com/kubeedge/kubeedge/edgemesh/pkg/resolver"
-)
-
-func Start() {
- //Initialize the resolvers
- r := &resolver.MyResolver{"http"}
- resolver.RegisterResolver(r)
- //Initialize the handlers
- //
- //
- config.GlobalDefinition = &model.GlobalCfg{}
- config.GlobalDefinition.Panel.Infra = "fake"
- opts := control.Options{
- Infra: config.GlobalDefinition.Panel.Infra,
- Address: config.GlobalDefinition.Panel.Settings["address"],
- }
- config.GlobalDefinition.Ssl = make(map[string]string)
-
- control.Init(opts)
- opt := registry.Options{}
- registry.DefaultServiceDiscoveryService = edgeregistry.NewServiceDiscovery(opt)
- myStrategy := meshconfig.Config.LBStrategy
- loadbalancer.InstallStrategy(myStrategy, func() loadbalancer.Strategy {
- switch myStrategy {
- case "RoundRobin":
- return &loadbalancer.RoundRobinStrategy{}
- case "Random":
- return &loadbalancer.RandomStrategy{}
- default:
- return &loadbalancer.RoundRobinStrategy{}
- }
- })
- //Start dns server
- go DnsStart()
- //Start server
- StartTCP()
-}
diff --git a/edgemesh/pkg/server/tcp.go b/edgemesh/pkg/server/tcp.go
deleted file mode 100644
index 0624723c2..000000000
--- a/edgemesh/pkg/server/tcp.go
+++ /dev/null
@@ -1,120 +0,0 @@
-package server
-
-import (
- "io"
- "io/ioutil"
- "net"
- "net/http"
-
- "github.com/go-chassis/go-chassis/core/common"
- "github.com/go-chassis/go-chassis/core/handler"
- "github.com/go-chassis/go-chassis/core/invocation"
- "k8s.io/klog"
-
- "github.com/kubeedge/kubeedge/edgemesh/pkg/resolver"
-)
-
-func StartTCP() {
- server, err := getIP()
- if err != nil {
- klog.Errorf("TCP server start error : %s", err)
- return
- }
-
- serverIP := server.String()
- // TODO Set as configurable @kadisi
- port := "8080"
- klog.Infof("start listening at %s:%s", serverIP, port)
- listener, err := net.Listen("tcp", serverIP+":"+port)
- if err != nil {
- klog.Errorf("failed to start TCP server with error:%v\n", err)
- return
- }
-
- for {
- conn, err := listener.Accept()
- if err != nil {
- klog.Errorf("failed to accept, err: %v\n", err)
- continue
- }
-
- go process(conn)
- }
-}
-
-func httpResponseToStr(resp *http.Response) string {
- respString := resp.Proto + " " + resp.Status + "\n"
- for key, values := range resp.Header {
- respString += key + ": "
- for _, v := range values {
- respString += v + ", "
- }
- respString = respString[0 : len(respString)-2]
- respString += "\n"
- }
- b, _ := ioutil.ReadAll(resp.Body)
- respString += "\n" + string(b)
- return respString
-}
-
-func process(conn net.Conn) {
- klog.Info("start receiving data...\n")
-
- buffer := make([]byte, 1024)
- d := make(chan []byte, 1024)
- s := make(chan interface{}, 1)
- restResponse := func(data *invocation.Response) error {
- if data.Err != nil {
- klog.Errorf("error in response:%v", data.Err)
- conn.Write([]byte(data.Err.Error()))
- return data.Err
- } else {
- if data.Result != nil {
- conn.Write([]byte(httpResponseToStr(data.Result.(*http.Response))))
- }
- }
- return nil
- }
- fakeResponse := func(data *invocation.Response) error {
- defer conn.Close()
- if data.Err != nil {
- klog.Errorf("error in response:%v", data.Err)
- conn.Write([]byte(data.Err.Error()))
- return data.Err
- } else {
- if data.Result != nil {
- conn.Write(data.Result.([]byte))
- }
- }
- return nil
- }
- invocationCallback := func(protocol string, invocation invocation.Invocation) {
- if invocation.Protocol == "rest" {
- c, err := handler.CreateChain(common.Consumer, protocol, handler.Loadbalance, handler.Transport)
- if err != nil {
- klog.Errorf("failed to create handlerchain:%v", err)
- }
- c.Next(&invocation, restResponse)
- } else {
- c, err := handler.CreateChain(common.Consumer, protocol)
- if err != nil {
- klog.Errorf("failed to create handlerchain:%v", err)
- }
- c.Next(&invocation, fakeResponse)
- }
- }
-
- //Start resolver
- go resolver.Resolve(d, s, invocationCallback)
- for {
- num, err := conn.Read(buffer)
- if err == nil {
- klog.Infof("buffer:\n%s\n", buffer)
- d <- buffer[:num]
- }
- if err == io.EOF {
- close(s)
- return
- }
- }
-}
diff --git a/edgemesh/pkg/server/tcp_test.go b/edgemesh/pkg/server/tcp_test.go
deleted file mode 100644
index a09e630ef..000000000
--- a/edgemesh/pkg/server/tcp_test.go
+++ /dev/null
@@ -1,249 +0,0 @@
-package server_test
-
-import (
- "bufio"
- "bytes"
- "fmt"
- "io/ioutil"
- "net"
- "net/http"
- "net/url"
- "strings"
- "testing"
- "time"
-
- "github.com/go-chassis/go-chassis/core/handler"
- "github.com/go-chassis/go-chassis/core/invocation"
- "k8s.io/klog"
-
- "github.com/kubeedge/kubeedge/edgemesh/pkg/resolver"
- "github.com/kubeedge/kubeedge/edgemesh/pkg/server"
-)
-
-type TestResolver struct {
- Name string
-}
-
-func httpMethods() (methods []string) {
- methods = []string{"GET", "HEAD", "POST", "OPTIONS", "PUT", "DELETE", "TRACE", "CONNECT"}
- return
-}
-
-func isHTTPRequest(s string) bool {
- methods := httpMethods()
- for _, method := range methods {
- if strings.HasPrefix(s, method) {
- return true
- }
- }
- return false
-}
-
-func testTransferHTTPRequest(req *http.Request) (invocation.Invocation, error) {
- clt := &http.Client{}
- u, err := url.Parse("http://127.0.0.1:9090")
- if err != nil {
- klog.Errorf("Parse new url error: %v\n", err)
- return invocation.Invocation{}, err
- }
- req.URL = u
- //clear RequestURI
- req.RequestURI = ""
-
- resp, err := clt.Do(req)
- if resp != nil {
- defer resp.Body.Close()
- }
- if err != nil {
- klog.Errorf("Resolve http request failed with error: %v\n", err)
- return invocation.Invocation{}, err
- }
- respBodyBytes, _ := ioutil.ReadAll(resp.Body)
- klog.Infof("resolve http resp body: %s\n", respBodyBytes)
- return invocation.Invocation{MicroServiceName: "http", Protocol: "rest", Args: resp}, nil
-}
-
-func (resolver *TestResolver) Resolve(data chan []byte, stop chan interface{}, invCallback func(string, invocation.Invocation)) (invocation.Invocation, bool) {
- content := ""
- protocol := ""
- for {
- select {
- case d := <-data:
- strData := string(d[:])
- if protocol == "" {
- if isHTTPRequest(strData) {
- protocol = "http"
- } else {
- return invocation.Invocation{}, false
- }
- }
- content += strData
- req, err := http.ReadRequest(bufio.NewReader(bytes.NewReader([]byte(content))))
- if err == nil {
- content = ""
- i, err := testTransferHTTPRequest(req)
- if err != nil {
- panic(err)
- }
- invCallback("http", i)
- }
- case <-stop:
- i := invocation.Invocation{MicroServiceName: resolver.Name, Args: content}
- invCallback(protocol, i)
- return i, true
- }
- klog.Infof("content: %s\n", content)
- }
-}
-
-type TestHandler struct{}
-
-//Handle
-func (h *TestHandler) Handle(chain *handler.Chain, inv *invocation.Invocation, cb invocation.ResponseCallBack) {
- r := &invocation.Response{
- Err: nil,
- }
- resp := "HTTP/1.1 200\ncontent-type: text/plain; charset=utf-8\ncontent-length: 20\ndate: Wed, 12 Jun 2019 09:28:08 GMT\n\n{\"name\": \"Jack\"}"
- r.Result = []uint8(resp)
- cb(r)
-}
-
-//Name
-func (h *TestHandler) Name() string {
- return "test"
-}
-func newTestHandler() handler.Handler {
- return &TestHandler{}
-}
-
-var serverStarted bool = false
-
-var httpServerStarted bool = false
-
-func StartTCPServer() {
- if serverStarted == false {
- //Initialize the resolvers
- r := &TestResolver{"http"}
- resolver.RegisterResolver(r)
- //Initialize the handlers
- handler.RegisterHandler("resolveHandler", newTestHandler)
- //Start server
- go server.StartTCP()
- serverStarted = true
- time.Sleep(1 * time.Second)
- }
-}
-
-func helloHTTP(w http.ResponseWriter, r *http.Request) {
- fmt.Fprintf(w, "Hello "+r.Method)
-}
-
-func StartHTTPServer() {
- if httpServerStarted == false {
- go func() {
- http.HandleFunc("/", helloHTTP)
- err := http.ListenAndServe(":9090", nil)
- if err != nil {
- klog.Errorf("ListenAndServe error: %v\n", err)
- }
- }()
- httpServerStarted = true
- time.Sleep(1 * time.Second)
- }
-}
-
-func handleTCPRead(conn net.Conn, done chan string) {
- buf := make([]byte, 1024)
- _, err := conn.Read(buf)
- if err != nil {
- klog.Infof("Error to read from TCP server: ", err)
- return
- }
- klog.Info("TCP server response: " + string(buf[:]))
-
- done <- "done"
-}
-
-func TestTCPRawBytes(t *testing.T) {
- //start TCP server if not started
- StartTCPServer()
-
- //start http server (:9090) to be transferred request
- go StartHTTPServer()
-
- //connect to server
- conn, err := net.Dial("tcp", "localhost:8080")
- if err != nil {
- t.Errorf("connect failed, err : %v\n", err.Error())
- return
- }
- defer conn.Close()
-
- //handle connection read
- done := make(chan string)
- go handleTCPRead(conn, done)
-
- //write raw bytes to TCP server
- testString := "GET / HTTP/1.1\nHost: 127.0.0.1:8080\n"
- _, err = conn.Write([]byte(testString))
- if err != nil {
- t.Errorf("write failed , err : %v\n", err)
- }
- time.Sleep(2 * time.Second)
- testString = "User-Agent: Go-http-client/1.1\nAccept-Encoding: gzip\n\n"
- _, err = conn.Write([]byte(testString))
- if err != nil {
- t.Errorf("write failed , err : %v\n", err)
- }
-
- <-done
- time.Sleep(3 * time.Second)
-}
-
-func TestResolveHTTP(t *testing.T) {
- //start TCP server if not started
- StartTCPServer()
-
- //start http server (:9090) to be transferred request
- StartHTTPServer()
-
- //new http client
- clt := &http.Client{}
-
- //do GET request
- req, err := http.NewRequest("GET", "http://127.0.0.1:8080", nil)
- if err != nil {
- t.Errorf("new http GET request error: %v\n", err)
- return
- }
- resp1, err := clt.Do(req)
- if resp1 != nil {
- defer resp1.Body.Close()
- }
- if err != nil {
- t.Errorf("do http GET request failed with error: %v\n", err)
- return
- }
- klog.Infof("GET response: %v\n", resp1)
-
- //do POST request
- data := url.Values{"Name": {"Mark"}, "Age": {"20"}}
- body := strings.NewReader(data.Encode())
- req, err = http.NewRequest("POST", "http://127.0.0.1:8080", body)
- if err != nil {
- t.Errorf("new http POST request error: %v\n", err)
- return
- }
- req.Header.Set("Content-Type", "application/x-www-form-urlencoded")
- resp2, err := clt.Do(req)
- if resp2 != nil {
- defer resp2.Body.Close()
- }
- if err != nil {
- t.Errorf("do http POST request failed with error: %v\n", err)
- return
- }
- klog.Infof("POST response: %v\n", resp2)
-
- time.Sleep(3 * time.Second)
-}