diff options
| author | liuzhiyi1993 <liuzhiyi@huawei.com> | 2020-02-13 05:56:21 +0800 |
|---|---|---|
| committer | liuzhiyi1993 <liuzhiyi@huawei.com> | 2020-03-31 22:30:02 +0800 |
| commit | 0dfe63446ad76e41a084ade9d3c393a2554def34 (patch) | |
| tree | f9857c1fd4ad3e9a00aa04effdb1c57bde6398b7 /edgemesh | |
| parent | Merge pull request #1567 from chendave/broken (diff) | |
| download | kubeedge-0dfe63446ad76e41a084ade9d3c393a2554def34.tar.gz | |
refactor edgemesh
Diffstat (limited to 'edgemesh')
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, ®istry.MicroServiceInstance{ + InstanceID: "", + ServiceID: name + "." + namespace, + HostName: "", + EndpointsMap: map[string]string{"rest": fmt.Sprintf("%s:%d", p.Status.HostIP, hostPort)}, + }) + } + } + + return microServiceInstances, nil +} + +// GetMicroServiceID get microServiceID +func (esd *EdgeServiceDiscovery) GetMicroServiceID(appID, microServiceName, version, env string) (string, error) { + return "", nil +} + +// GetMicroServiceInstances return instances +func (esd *EdgeServiceDiscovery) GetMicroServiceInstances(consumerID, providerID string) ([]*registry.MicroServiceInstance, error) { + return nil, nil +} + +// GetMicroService return service +func (esd *EdgeServiceDiscovery) GetMicroService(microServiceID string) (*registry.MicroService, error) { + return nil, nil +} + +// AutoSync updating the cache manager +func (esd *EdgeServiceDiscovery) AutoSync() {} + +// Close close all websocket connection +func (esd *EdgeServiceDiscovery) Close() error { return nil } + +// parseServiceUrl parses serviceUrl to ${service_name}.${namespace}.svc.${cluster}:${port}, keeps with k8s service +func parseServiceUrl(serviceUrl string) (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, ®istry.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) -} |
