diff --git a/third_party/meshnet/daemon/cni/cni.go b/third_party/meshnet/daemon/cni/cni.go index b3611855..f0df080b 100644 --- a/third_party/meshnet/daemon/cni/cni.go +++ b/third_party/meshnet/daemon/cni/cni.go @@ -1,3 +1,4 @@ +// Package cni handles CNI configuration file installation and cleanup for meshnet. package cni import ( diff --git a/third_party/meshnet/daemon/grpcwire/grpcwire.go b/third_party/meshnet/daemon/grpcwire/grpcwire.go index e9d2181e..870bac40 100644 --- a/third_party/meshnet/daemon/grpcwire/grpcwire.go +++ b/third_party/meshnet/daemon/grpcwire/grpcwire.go @@ -1,27 +1,24 @@ +// Package grpcwire provides gRPC overlay wire creation, TAP interface management, +// packet multiplexing, and CRD reconciliation for meshnet daemon. package grpcwire import ( - "context" "fmt" "io" - "net" - "strings" + "os" "sync" - "github.com/google/gopacket" - "github.com/google/gopacket/pcap" + "github.com/containernetworking/plugins/pkg/ns" "github.com/openconfig/gnmi/errlist" - koko "github.com/redhat-nfvpe/koko/api" log "github.com/sirupsen/logrus" - "google.golang.org/grpc" - "google.golang.org/grpc/credentials/insecure" + "github.com/vishvananda/netlink" mpb "github.com/openconfig/kne/third_party/meshnet/daemon/proto/meshnet/v1beta1" - "github.com/openconfig/kne/third_party/meshnet/utils/wireutil" ) var grpcOvrlyLogger *log.Entry = nil +// InitLogger initializes logrus logging for the gRPC overlay daemon. func InitLogger() { grpcOvrlyLogger = log.WithFields(log.Fields{"daemon": "meshnetd", "overlay": "gRPC"}) } @@ -32,16 +29,14 @@ type intfIndex struct { } /* - In a given node a veth-pair connects a pod with the meshnet daemon hosted in the node. This meshnet + In a given node a TAP interface connects a pod with the meshnet daemon hosted in the node. This meshnet -daemon provides the grpc-wire service to connect the local pod with the remote pod over grpc. The node -end of the veth-pair must have unique name with in the node. A node can have multiple pods. So there -will be multiple veth-pairs for connecting multiple nodes to meshnet daemon and each of them (the node end) must have unique -names. IntfIndex provides the sequentially increasing number which makes the name unique when added as -suffix to the name. +daemon provides the grpc-wire service to connect the local pod with the remote pod over grpc. IntfIndex provides +the sequentially increasing number which makes the name unique when added as suffix to the name. */ var indexGen intfIndex +// NextIndex generates a node-wide, monotonically increasing unique wire ID for TAP device handle indexing. func NextIndex() int64 { indexGen.mu.Lock() defer indexGen.mu.Unlock() @@ -49,7 +44,6 @@ func NextIndex() int64 { return indexGen.currId } -/*+++tbf: These constants has no utility other that helping in debugging. These can be removed later. */ type grpcWireOriginator int func (g grpcWireOriginator) String() string { @@ -98,6 +92,7 @@ type linkKey struct { linkUID int } +// CreateGWire constructs a new GRPCWire struct from the provided wire definition. func CreateGWire(locIfIndex int, locIfNm string, stopC chan struct{}, wireDef *mpb.WireDef) *GRPCWire { return &GRPCWire{ @@ -189,14 +184,9 @@ func WireDownByUID(namespace string, linkUID int) error { return nil } -// ------------------------------------------------------------------------------------------------- -func AddWireInMemNDataStore(wire *GRPCWire, handle *pcap.Handle) int { +// AddWireInMemNDataStore populates the active wire map and updates K8s status store. +func AddWireInMemNDataStore(wire *GRPCWire, handle *os.File) int { /* Populate the active wire map and returns the number of currently added active wires. */ - - /* if this wire is already present in the map then it will be overwritten. - It seems to be ok to overwrite. Think more in what situation this may - not be the desired behavior and we need to throw an error. */ - wires.AddInMemNDataStore(wire, handle) return len(wires.wires) } @@ -242,7 +232,7 @@ func DeletePodWires(namespace string, podName string) error { func RemoveWireAcrosAll(wire *GRPCWire, inMem bool) error { if wire == nil { - grpcOvrlyLogger.Infof("[WIRE-DELETE]:Null wire. This ware is already removed") + grpcOvrlyLogger.Infof("[WIRE-DELETE]:Null wire. This wire is already removed") return nil } @@ -252,19 +242,24 @@ func RemoveWireAcrosAll(wire *GRPCWire, inMem bool) error { } wire.IsReady = false - /* Remove the veth from the node */ - intf, err := net.InterfaceByIndex(int(wire.LocalNodeIfaceID)) - if err != nil { - grpcOvrlyLogger.Infof("[WIRE-DELETE]:Interface index %d for wire %d, is already cleaned up.", wire.LocalNodeIfaceID, wire.UID) - } else { - myVeth := koko.VEth{} - myVeth.LinkName = intf.Name - if err = myVeth.RemoveVethLink(); err != nil { - return fmt.Errorf("[WIRE-DELETE]:failed to remove veth link: %w", err) - } + // Close the TAP file handle if open + if handle, ok := wires.GetHandle(wire.LocalNodeIfaceID); ok && handle != nil { + _ = handle.Close() + } + + // Remove the TAP link from the container netns if present + podNs, err := ns.GetNS(wire.LocalPodNetNS) + if err == nil { + _ = podNs.Do(func(_ ns.NetNS) error { + if link, err := netlink.LinkByName(wire.LocalPodIfaceName); err == nil { + return netlink.LinkDel(link) + } + return nil + }) + podNs.Close() } - // clean up im-memory wire-map + // clean up in-memory wire-map if inMem { wires.AtomicDelete(wire) // Deleting the wire from in-memory data } @@ -277,112 +272,79 @@ func RemoveWireAcrosAll(wire *GRPCWire, inMem bool) error { // ----------------------------------------------------------------------------------------------------------- // Generate the name of the interface to be placed on the node func GenNodeIfaceName(podName string, podIfaceName string) (string, error) { - // Linux has issue if interface name is too long. Generate a smaller name. - // In recent kernel versions this is defined by IFNAMSIZ to be 16 bytes, so 15 user-visible bytes - // (assuming it includes a trailing null). IFNAMSIZ is used in defining struct net_device's name. - // The name must not contain / or any whitespace characters - // - //TODO: This method needs to be robust. It monotonically increases the index and never - // decreases it, even if the interfaces are deleted. So far this will work for accumulated - // 1K interfaces per node under the current naming scheme. This is too small. - // Using 14 digit random number and checking if any interface with generated name exists and if - // exists then generate another random number (try 3 times before giving up). This will make it robust. - // This reduces the readability and correlation between the “pod-interface” and corresponding - // “node-interface”, for example eth1host1-<3-digit-index> will become "12345678901234". id := NextIndex() - ifaceName := fmt.Sprintf("%.5s%.5s-%04d", podName, podIfaceName, id) - return ifaceName, nil } -// ----------------------------------------------------------------------------------------------------------- +// RecvFrmLocalPodThread reads packets from the local TAP interface and forwards them over the gRPC stream. func RecvFrmLocalPodThread(wire *GRPCWire, locIfNm string) error { - - defaultPort := wireutil.GRPCDefaultPort - pktBuffSz := int32(1024 * 64 * 10) //keep buffer for MAX 10 64K frames - - url := strings.TrimSpace(fmt.Sprintf("%s:%d", wire.PeerNodeIP, defaultPort)) - /* Utilizing google gopacket for polling for packets from the node. This seems to be the - simplest way to get all packets. - As an alternative to google gopacket(pcap), a socket based implementation is possible. - Not sure if socket based implementation can bring any advantage or not. - - Near term will replace pcap by socket. - */ - - // in some rare cases by the time the thread starts K8S may decide to move the pod somewhere else. - // in that case the local interfaced will be cleaned up asynchronously. Detect the situation and return. - _, err := net.InterfaceByName(locIfNm) + tapFile, err := GetHostIntfHndl(wire.LocalNodeIfaceID) if err != nil { - grpcOvrlyLogger.Errorf("[Packet Receive thread]For pod %s failed to retrieve interface %s/%d. error: %v", wire.LocalPodName, wire.LocalNodeIfaceName, wire.LocalNodeIfaceID, err) + grpcOvrlyLogger.Errorf("[Packet Receive thread] For pod %s failed to retrieve TAP handle for interface %s/%d. error: %v", wire.LocalPodName, wire.LocalNodeIfaceName, wire.LocalNodeIfaceID, err) return err } - rdHandl, err := pcap.OpenLive(wire.LocalNodeIfaceName, pktBuffSz, true, pcap.BlockForever) - if err != nil { - // let the caller handle the error - grpcOvrlyLogger.Errorf("Receive Thread for local pod failed to open interface: %s/%d, PCAP ERROR: %v", wire.LocalNodeIfaceName, wire.LocalNodeIfaceID, err) - return err - } - defer rdHandl.Close() + nodeStream := streamMgr.GetOrCreateStream(wire.TopoNamespace, wire.PeerNodeIP) + defer streamMgr.ReleaseStream(wire.TopoNamespace, wire.PeerNodeIP) - err = rdHandl.SetDirection(pcap.Direction(pcap.DirectionIn)) - if err != nil { - // let the caller handle the error - grpcOvrlyLogger.Errorf("Receive Thread for local pod failed to set up capture direction: %s/%d, PCAP ERROR: %v", wire.LocalNodeIfaceName, wire.LocalNodeIfaceID, err) - return err + buf := make([]byte, 65535) + type readResult struct { + n int + err error } + readChan := make(chan readResult, 1) + go func() { + for { + n, err := tapFile.Read(buf) + readChan <- readResult{n: n, err: err} + if err != nil { + return + } + } + }() - remote, err := grpc.Dial(url, grpc.WithTransportCredentials(insecure.NewCredentials())) - if err != nil { - grpcOvrlyLogger.Infof("RecvFrmLocalPodThread:Failed to connect to remote %s/%d", url, wire.LocalNodeIfaceID) - return err - } - defer remote.Close() - - ctx, cancel := context.WithCancel(context.Background()) - defer cancel() - - source := gopacket.NewPacketSource(rdHandl, rdHandl.LinkType()) - wireClient := mpb.NewWireProtocolClient(remote) - - in := source.Packets() - var packet gopacket.Packet for { select { case <-wire.StopC: grpcOvrlyLogger.Infof("RecvFrmLocalPodThread: closing connection with remote peer-iface@peer-node-ip: %d@%s/%d from %s@%s", wire.WireIfaceIDOnPeerNode, wire.PeerNodeIP, wire.LocalNodeIfaceID, wire.LocalPodName, wire.LocalPodIfaceName) return io.EOF - case packet = <-in: - data := packet.Data() + case res := <-readChan: + if res.err != nil { + select { + case <-wire.StopC: + return io.EOF + default: + grpcOvrlyLogger.Errorf("RecvFrmLocalPodThread: error reading from TAP interface %s: %v", locIfNm, res.err) + return res.err + } + } + n := res.n + if n <= 0 { + continue + } + + frame := make([]byte, n) + copy(frame, buf[:n]) + + if !wire.IsReady || wire.WireIfaceIDOnPeerNode <= 0 { + // Remote peer handshake is still in progress; skip sending to unassigned wire ID 0 + continue + } + payload := &mpb.Packet{ RemotIntfId: wire.WireIfaceIDOnPeerNode, - Frame: data, + Frame: frame, } - /*+++TODO: Ethernet has a minimum frame size of 64 bytes, comprising an 18-byte header and a payload of 46 bytes. - It also has a maximum frame size of 1518 bytes, in which case the payload is 1500 bytes. - This logic needs to be better, take the interface MTU not hardcoded value of 1518. - This is a very unusual condition to receive an packet from the pod with size > MTU. This can only happens if - things gets really messed up. */ - if len(data) > 1518 { + if n > 1518 { pktType := DecodeFrame(payload.Frame) - grpcOvrlyLogger.Infof("RecvFrmLocalPodThread: unusually large packet received from local pod (may be GRO enabled). size: %d, pkt:%s", len(data), pktType) - /* When Generic Receive Offload (GRO) is enabled then containers can send packets larger than MTU size packet. Do not drop these - packets, deliver it to the receiving container to process. - */ - //continue + grpcOvrlyLogger.Infof("RecvFrmLocalPodThread: unusually large packet received from local pod (may be GRO enabled). size: %d, pkt:%s", n, pktType) } - ok, err := wireClient.SendToOnce(ctx, payload) - if err != nil || !ok.Response { - grpcOvrlyLogger.Infof("RecvFrmLocalPodThread: Could not deliver pkt %s@%s@%s. Peer not ready, remote iface id %d. err=%v", - wire.LocalPodName, wire.LocalPodIfaceName, wire.LocalNodeIfaceName, wire.WireIfaceIDOnPeerNode, err) - /* we generate information and continue. As the above errors will happen when the remote end is not yet ready. - It will eventually get ready and if it can't then someone else will stop this thread. - */ + if !nodeStream.Send(payload) { + grpcOvrlyLogger.Debugf("RecvFrmLocalPodThread: Could not queue packet over stream %s@%s (queue full)", wire.LocalPodName, wire.LocalNodeIfaceName) } } } diff --git a/third_party/meshnet/daemon/grpcwire/gwire_map.go b/third_party/meshnet/daemon/grpcwire/gwire_map.go index b58fbc33..f583e135 100644 --- a/third_party/meshnet/daemon/grpcwire/gwire_map.go +++ b/third_party/meshnet/daemon/grpcwire/gwire_map.go @@ -2,15 +2,14 @@ package grpcwire import ( "fmt" + "os" "sync" - - "github.com/google/gopacket/pcap" ) type wireMap struct { mu sync.Mutex wires map[linkKey]*GRPCWire - handles map[int64]*pcap.Handle + handles map[int64]*os.File } func (w *wireMap) GetWire(namespace string, linkUID int) (*GRPCWire, bool) { @@ -23,14 +22,14 @@ func (w *wireMap) GetWire(namespace string, linkUID int) (*GRPCWire, bool) { return wire, ok } -func (w *wireMap) GetHandle(key int64) (*pcap.Handle, bool) { +func (w *wireMap) GetHandle(key int64) (*os.File, bool) { w.mu.Lock() defer w.mu.Unlock() handle, ok := w.handles[key] return handle, ok } -func (w *wireMap) AddInMem(wire *GRPCWire, handle *pcap.Handle) error { +func (w *wireMap) AddInMem(wire *GRPCWire, handle *os.File) error { w.mu.Lock() defer w.mu.Unlock() w.wires[linkKey{ @@ -42,17 +41,16 @@ func (w *wireMap) AddInMem(wire *GRPCWire, handle *pcap.Handle) error { return nil } -func (w *wireMap) AddInMemNDataStore(wire *GRPCWire, handle *pcap.Handle) error { +func (w *wireMap) AddInMemNDataStore(wire *GRPCWire, handle *os.File) error { w.mu.Lock() - defer w.mu.Unlock() w.wires[linkKey{ namespace: wire.LocalPodNetNS, linkUID: wire.UID, }] = wire - - wire.K8sStoreGWire() - w.handles[wire.LocalNodeIfaceID] = handle + w.mu.Unlock() + + go wire.K8sStoreGWire() return nil } @@ -87,11 +85,8 @@ func (w *wireMap) DeleteWoLock(wire *GRPCWire) error { * trigger is received. This situation occurs when both the host triggers wire creation almost simultaneously. */ var wires = &wireMap{ - wires: map[linkKey]*GRPCWire{}, - /* Used when a packet is received, then we know the id of the interface to which the packet to be delivered. - This map take interface-id as key and returns the corresponding handle for delivering the packet. - map[interface-id]->handle */ - handles: map[int64]*pcap.Handle{}, + wires: map[linkKey]*GRPCWire{}, + handles: map[int64]*os.File{}, } // FindWiresByPod returns a list of wires matching the namespace and pod. @@ -113,7 +108,6 @@ func GetWiresByPod(namespace string, podName string) ([]*GRPCWire, bool) { func ExtractOneWireByPod(namespace string, podName string) (*GRPCWire, bool) { wires.mu.Lock() defer wires.mu.Unlock() - //var rWires *GRPCWire for _, wire := range wires.wires { if wire.LocalPodName == podName && wire.TopoNamespace == namespace { @@ -123,7 +117,7 @@ func ExtractOneWireByPod(namespace string, podName string) (*GRPCWire, bool) { linkUID: wire.UID, }) - // also clean up the pcap handle for the wire that is extracted from the wire-map + // also clean up the handle for the wire that is extracted from the wire-map delete(wires.handles, wire.LocalNodeIfaceID) return wire, true } @@ -131,7 +125,7 @@ func ExtractOneWireByPod(namespace string, podName string) (*GRPCWire, bool) { return nil, true // no wire found is not a failure, so return true } -func GetHostIntfHndl(intfID int64) (*pcap.Handle, error) { +func GetHostIntfHndl(intfID int64) (*os.File, error) { val, ok := wires.GetHandle(intfID) if ok { diff --git a/third_party/meshnet/daemon/grpcwire/gwire_recon.go b/third_party/meshnet/daemon/grpcwire/gwire_recon.go index a97cb8bd..592922d2 100644 --- a/third_party/meshnet/daemon/grpcwire/gwire_recon.go +++ b/third_party/meshnet/daemon/grpcwire/gwire_recon.go @@ -3,13 +3,13 @@ package grpcwire import ( "context" "fmt" - "net" "os" "reflect" + "sync" - "github.com/google/gopacket/pcap" grpcwirev1 "github.com/openconfig/kne/third_party/meshnet/api/types/v1beta1" mpb "github.com/openconfig/kne/third_party/meshnet/daemon/proto/meshnet/v1beta1" + "github.com/openconfig/kne/third_party/meshnet/utils/wireutil" log "github.com/sirupsen/logrus" "k8s.io/apimachinery/pkg/api/errors" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" @@ -34,7 +34,7 @@ const ( kGrpcWireItems = "grpcWireItems" // json name of GWireKItems of gwire_type, +++TBD: can we make it dynamic ) -// ----------------------------------------------------------------------------------------------------------- +// SetGWireClient initializes the dynamic K8s client for gRPC wire CRD management. func SetGWireClient(gClient *dynamic.DynamicClient) { // identifier of grpc wire object in k8s apis gWClient.gvr = schema.GroupVersionResource{ @@ -45,12 +45,12 @@ func SetGWireClient(gClient *dynamic.DynamicClient) { gWClient.di = gClient.Resource(gWClient.gvr) } -// ----------------------------------------------------------------------------------------------------------- +// SetGWireClientInterface sets the K8s dynamic resource interface (used for unit testing). func SetGWireClientInterface(gClient dynamic.NamespaceableResourceInterface) { gWClient.di = gClient } -// ------------------------------------------------------------------------------------------------------------ +// GetWireObjListUS lists unstructured GWireKObj resources for a specified node. func (gc GWireClient) GetWireObjListUS(ctx context.Context, ndName string) (*unstructured.UnstructuredList, error) { return gc.di.Namespace("").List(ctx, metav1.ListOptions{ TypeMeta: metav1.TypeMeta{ @@ -62,19 +62,17 @@ func (gc GWireClient) GetWireObjListUS(ctx context.Context, ndName string) (*uns }) } -// ------------------------------------------------------------------------------------------------------------ +// CreatWireObj creates a new unstructured GWireKObj resource in K8s. func (gc GWireClient) CreatWireObj(ctx context.Context, nSpace string, uWbj map[string]interface{}) (*unstructured.Unstructured, error) { return gc.di.Namespace(nSpace).Create(ctx, &unstructured.Unstructured{Object: uWbj}, metav1.CreateOptions{}) } -// ------------------------------------------------------------------------------------------------------------ +// UpdateWireObj updates an existing unstructured GWireKObj resource in K8s. func (gc GWireClient) UpdateWireObj(ctx context.Context, nSpace string, wObjsOnNd *unstructured.Unstructured) (*unstructured.Unstructured, error) { return gc.di.Namespace(nSpace).Update(ctx, wObjsOnNd, metav1.UpdateOptions{}) - } -//------------------------------------------------------------------------------------------------------------ - +// GetWireObjGrpUS retrieves the GWireKObj for a given node and status. func (gc GWireClient) GetWireObjGrpUS(ctx context.Context, wStatus *grpcwirev1.GWireStatus) (*unstructured.Unstructured, error) { return gc.di.Namespace(wStatus.TopoNamespace).Get(ctx, wStatus.LocalNodeName, metav1.GetOptions{}) } @@ -102,22 +100,153 @@ func CreateWireStatus(wire *GRPCWire, nodeName string) *grpcwirev1.GWireStatus { } +type wireStatusUpdate struct { + wire *GRPCWire + nodeName string + flushDone chan struct{} +} + +var ( + statusQueueChan = make(chan wireStatusUpdate, 10000) + statusQueueOnce sync.Once +) + +func startStatusQueueWorker() { + statusQueueOnce.Do(func() { + go func() { + for { + first, ok := <-statusQueueChan + if !ok { + return + } + + var flushSignal chan struct{} + var updates []wireStatusUpdate + + if first.flushDone != nil { + flushSignal = first.flushDone + } else { + updates = append(updates, first) + } + + drain := true + for drain && len(updates) < 200 { + select { + case u, ok := <-statusQueueChan: + if !ok { + drain = false + } else if u.flushDone != nil { + flushSignal = u.flushDone + drain = false + } else { + updates = append(updates, u) + } + default: + drain = false + } + } + + if len(updates) > 0 { + if err := updateGRPCWireStatusBatch(context.Background(), updates); err != nil { + grpcOvrlyLogger.Warnf("K8sStatusWorker: failed to update batch of %d wire statuses: %v", len(updates), err) + } + } + + if flushSignal != nil { + close(flushSignal) + } + } + }() + }) +} + +// FlushK8sStatusQueue blocks until all currently queued status updates have been processed by the worker. +func FlushK8sStatusQueue() { + startStatusQueueWorker() + done := make(chan struct{}) + statusQueueChan <- wireStatusUpdate{flushDone: done} + <-done +} + +func updateGRPCWireStatusBatch(ctx context.Context, updates []wireStatusUpdate) error { + if len(updates) == 0 { + return nil + } + + nodeName := updates[0].nodeName + topoNs := updates[0].wire.TopoNamespace + + return retry.RetryOnConflict(retry.DefaultRetry, func() error { + wStatusFirst := CreateWireStatus(updates[0].wire, nodeName) + wObjsOnNd, err := gWClient.GetWireObjGrpUS(ctx, wStatusFirst) + if err != nil { + if errors.IsNotFound(err) { + err = CreateGWireStatInDS(ctx, wStatusFirst) + if err != nil { + return err + } + wObjsOnNd, err = gWClient.GetWireObjGrpUS(ctx, wStatusFirst) + if err != nil { + return err + } + } else { + return err + } + } + + gwireItems, found, err := unstructured.NestedSlice(wObjsOnNd.Object, kStatus, kGrpcWireItems) + if err != nil || !found || gwireItems == nil { + gwireItems = []interface{}{} + } + + existingMap := make(map[string]int) + for i, item := range gwireItems { + if m, ok := item.(map[string]interface{}); ok { + key := fmt.Sprintf("%v@%v", m["local_pod_name"], m["local_pod_iface_name"]) + existingMap[key] = i + } + } + + for _, u := range updates { + ws := CreateWireStatus(u.wire, u.nodeName) + unstrucObj, err := runtime.DefaultUnstructuredConverter.ToUnstructured(&ws) + if err != nil { + continue + } + + key := fmt.Sprintf("%v@%v", ws.LocalPodName, ws.LocalPodIfaceName) + if idx, exists := existingMap[key]; exists { + gwireItems[idx] = unstrucObj + } else { + gwireItems = append(gwireItems, unstrucObj) + existingMap[key] = len(gwireItems) - 1 + } + } + + if err := unstructured.SetNestedField(wObjsOnNd.Object, gwireItems, kStatus, kGrpcWireItems); err != nil { + return err + } + + _, err = gWClient.UpdateWireObj(ctx, topoNs, wObjsOnNd) + return err + }) +} + // ----------------------------------------------------------------------------------------------------------- // K8sStoreGWire writes grpc wire info 'wire' for a specific topology namespace (wire.TopoNamespace) into k8s // data-store for the current node. It calls updateGRPCWireStatus() to serve the purpose func (wire *GRPCWire) K8sStoreGWire() error { + startStatusQueueWorker() nodeName, err := findNodeName() if err != nil { grpcOvrlyLogger.Errorf("K8sStoreGWire: could not get node name: %v", err) return err } - ctx := context.Background() - ws := CreateWireStatus(wire, nodeName) - err = updateGRPCWireStatus(ctx, ws) - - if err != nil { - grpcOvrlyLogger.Errorf("K8sStoreGWire: Failed to set status for node %s: %v", nodeName, err) + select { + case statusQueueChan <- wireStatusUpdate{wire: wire, nodeName: nodeName}: + default: + grpcOvrlyLogger.Warnf("K8sStoreGWire: status queue full, dropping async status update for %s@%s", wire.LocalPodName, wire.LocalPodIfaceName) } return nil } @@ -422,24 +551,18 @@ func reCreateGWire(wStatus grpcwirev1.GWireStatus, _ context.Context) error { // ----------------------------------------------------------------------------------------------------------- // Recreate the wire in-memory wire-map and start the pod to daemon packet receive thread for this wire. func reconLocalGRPCWire(wireDef *mpb.WireDef) error { - locInf, err := net.InterfaceByName(wireDef.WireIfNameOnLocalNode) - if err != nil { - grpcOvrlyLogger.Errorf("[RECONCILE:LOCAL-END]For pod %s failed to retrieve interface ID for interface %v. error:%v", wireDef.LocalPodName, wireDef.WireIfNameOnLocalNode, err) - return err - } - - //Using google gopacket for packet receive. An alternative could be using socket. Not sure it it provides any advantage over gopacket. - wrHandle, err := pcap.OpenLive(wireDef.WireIfNameOnLocalNode, 65365, true, pcap.BlockForever) + tapFile, err := wireutil.CreateOrAttachTAP(wireDef.LocalPodNetNs, wireDef.IntfNameInPod, wireDef.LocalPodIp) if err != nil { - grpcOvrlyLogger.Errorf("[RECONCILE:LOCAL-END]Could not open interface for send/recv packets for containers local iface id %d. error:%v", locInf.Index, err) + grpcOvrlyLogger.Errorf("[RECONCILE:LOCAL-END] For pod %s failed to create/attach TAP interface %s in netns %s: %v", + wireDef.LocalPodName, wireDef.IntfNameInPod, wireDef.LocalPodNetNs, err) return err } - aWire := CreateGWire(locInf.Index, wireDef.WireIfNameOnLocalNode, make(chan struct{}), wireDef) + wireID := NextIndex() + aWire := CreateGWire(int(wireID), wireDef.IntfNameInPod, make(chan struct{}), wireDef) aWire.IsReady = true // reconciling, so add only in memory - wires.AddInMem(aWire, wrHandle) + wires.AddInMem(aWire, tapFile) - // TODO: handle error here go RecvFrmLocalPodThread(aWire, aWire.LocalNodeIfaceName) return nil diff --git a/third_party/meshnet/daemon/grpcwire/gwire_recon_test.go b/third_party/meshnet/daemon/grpcwire/gwire_recon_test.go index bdbc6bff..6ee631cb 100644 --- a/third_party/meshnet/daemon/grpcwire/gwire_recon_test.go +++ b/third_party/meshnet/daemon/grpcwire/gwire_recon_test.go @@ -272,6 +272,7 @@ func TestK8sStoreGWire(t *testing.T) { if err != nil { t.Fatalf("could not add gwire status into k8s data-store") } + FlushK8sStatusQueue() storedWStatus = append(storedWStatus, *CreateWireStatus(tc.store, nodeName)) wObjsOnNd, err := cs.Namespace(tc.store.TopoNamespace).Get(context.Background(), nodeName, metav1.GetOptions{}) if err != nil { diff --git a/third_party/meshnet/daemon/grpcwire/gwire_rpc_handlers.go b/third_party/meshnet/daemon/grpcwire/gwire_rpc_handlers.go index a5d91ebd..01863394 100644 --- a/third_party/meshnet/daemon/grpcwire/gwire_rpc_handlers.go +++ b/third_party/meshnet/daemon/grpcwire/gwire_rpc_handlers.go @@ -2,60 +2,38 @@ package grpcwire import ( "context" - "fmt" - "net" log "github.com/sirupsen/logrus" - "github.com/containernetworking/plugins/pkg/ns" - "github.com/google/gopacket/pcap" mpb "github.com/openconfig/kne/third_party/meshnet/daemon/proto/meshnet/v1beta1" "github.com/openconfig/kne/third_party/meshnet/utils/wireutil" - koko "github.com/redhat-nfvpe/koko/api" ) func CreateGRPCWireLocal(ctx context.Context, wireDef *mpb.WireDef) (*mpb.BoolResponse, error) { - locInf, err := net.InterfaceByName(wireDef.WireIfNameOnLocalNode) + tapFile, err := wireutil.CreateOrAttachTAP(wireDef.LocalPodNetNs, wireDef.IntfNameInPod, wireDef.LocalPodIp) if err != nil { log.WithFields(log.Fields{ "daemon": "meshnetd", "overlay": "gRPC", - }).Errorf("[ADD-WIRE:LOCAL-END]For pod %s failed to retrieve interface ID for interface %v. error:%v", wireDef.LocalPodName, wireDef.WireIfNameOnLocalNode, err) + }).Errorf("[ADD-WIRE:LOCAL-END] For pod %s failed to create/attach TAP interface %s in netns %s: %v", + wireDef.LocalPodName, wireDef.IntfNameInPod, wireDef.LocalPodNetNs, err) return &mpb.BoolResponse{Response: false}, err } - // update tx checksumming to off - err = wireutil.SetTxChecksumOff(wireDef.IntfNameInPod, wireDef.LocalPodNetNs) - if err != nil { - log.Errorf("Error in setting tx checksum-off on interface %s, ns %s, pod %s: %v", wireDef.IntfNameInPod, wireDef.LocalPodNetNs, wireDef.LocalPodName, err) - // generate error and continue - } else { - log.Infof("Setting tx checksum-off on interface %s, pod %s is successful", wireDef.IntfNameInPod, wireDef.LocalPodName) - } - - //Using google gopacket for packet receive. An alternative could be using socket. Not sure it it provides any advantage over gopacket. - wrHandle, err := pcap.OpenLive(wireDef.WireIfNameOnLocalNode, 65365, true, pcap.BlockForever) - if err != nil { - log.WithFields(log.Fields{ - "daemon": "meshnetd", - "overlay": "gRPC", - }).Errorf("[ADD-WIRE:LOCAL-END]Could not open interface for send/recv packets for containers local iface id %d. error:%v", locInf.Index, err) - return &mpb.BoolResponse{Response: false}, err - } - - aWire := CreateGWire(locInf.Index, wireDef.WireIfNameOnLocalNode, make(chan struct{}), wireDef) + wireID := NextIndex() + aWire := CreateGWire(int(wireID), wireDef.IntfNameInPod, make(chan struct{}), wireDef) aWire.IsReady = false aWire.Originator = HOST_CREATED_WIRE aWire.OriginatorIP = "unknown" // Add the newly created wire in the in memory wire-map and k8S data store - AddWireInMemNDataStore(aWire, wrHandle) + AddWireInMemNDataStore(aWire, tapFile) log.WithFields(log.Fields{ "daemon": "meshnetd", "overlay": "gRPC", - }).Infof("[ADD-WIRE:LOCAL-END]For pod %s@%s, node iface id %d starting the local packet receive thread", wireDef.LocalPodName, wireDef.IntfNameInPod, locInf.Index) - // TODO: handle error here + }).Infof("[ADD-WIRE:LOCAL-END] For pod %s@%s, wire id %d starting local packet receive thread", wireDef.LocalPodName, wireDef.IntfNameInPod, wireID) + go RecvFrmLocalPodThread(aWire, aWire.LocalNodeIfaceName) return &mpb.BoolResponse{Response: true}, nil @@ -67,8 +45,6 @@ func CreateGRPCWireLocal(ctx context.Context, wireDef *mpb.WireDef) (*mpb.BoolRe // a pod from node A to node B dynamically func CreateUpdateGRPCWireRemoteTriggered(wireDef *mpb.WireDef, stopC chan struct{}) (*GRPCWire, error) { - var err error - // If this wire is already created, then only update the already created wire properties like stopC. // This can happen due to a race between the local and remote peer. // This can also happen when a pod in one end of the wire is deleted and created again. @@ -79,73 +55,21 @@ func CreateUpdateGRPCWireRemoteTriggered(wireDef *mpb.WireDef, stopC chan struct return grpcWire, nil } - outIfNm, err := GenNodeIfaceName(wireDef.LocalPodName, wireDef.IntfNameInPod) - if err != nil { - return nil, fmt.Errorf("[ADD-WIRE:REMOTE-END] could not get current network namespace: %v", err) - } - - currNs, err := ns.GetCurrentNS() + tapFile, err := wireutil.CreateOrAttachTAP(wireDef.LocalPodNetNs, wireDef.IntfNameInPod, wireDef.LocalPodIp) if err != nil { - return nil, fmt.Errorf("[ADD-WIRE:REMOTE-END] could not get current network namespace: %v", err) - } - - /* Create the veth to connect the pod with the meshnet daemon running on the node */ - hostEndVeth := koko.VEth{ - NsName: currNs.Path(), - LinkName: outIfNm, - } - - inIfNm := wireDef.IntfNameInPod - inContainerVeth := koko.VEth{ - NsName: wireDef.LocalPodNetNs, - LinkName: inIfNm, + grpcOvrlyLogger.Errorf("[ADD-WIRE:REMOTE-END] Error creating/attaching TAP interface %s in netns %s: %v", + wireDef.IntfNameInPod, wireDef.LocalPodNetNs, err) + return nil, err } - if wireDef.LocalPodIp != "" { - ipAddr, ipSubnet, err := net.ParseCIDR(wireDef.LocalPodIp) - if err != nil { - return nil, fmt.Errorf("failed to create remote end of GRPC wire(%s@%s), failed to parse CIDR %s: %w", - inIfNm, wireDef.LocalPodName, wireDef.LocalPodIp, err) - } - inContainerVeth.IPAddr = []net.IPNet{{ - IP: ipAddr, - Mask: ipSubnet.Mask, - }} - } + wireID := NextIndex() + grpcOvrlyLogger.Infof("[ADD-WIRE:REMOTE-END] Trigger from %s:%d : Successfully created/attached TAP interface %s@%s (wire id %d).", + wireDef.PeerNodeIp, wireDef.WireIfIdOnPeerNode, wireDef.IntfNameInPod, wireDef.LocalPodName, wireID) - if err = koko.MakeVeth(inContainerVeth, hostEndVeth); err != nil { - grpcOvrlyLogger.Errorf("[ADD-WIRE:REMOTE-END] Error creating vEth pair (in:%s <--> out:%s). Error-> %s", inIfNm, outIfNm, err) - return nil, err - } - if err := wireutil.SetTxChecksumOff(inContainerVeth.LinkName, inContainerVeth.NsName); err != nil { - grpcOvrlyLogger.Errorf("Error in setting tx checksum-off on interface %s, pod %s: %v", inContainerVeth.LinkName, wireDef.LocalPodName, err) - // not returning - } - locIface, err := net.InterfaceByName(hostEndVeth.LinkName) - if err != nil { - // let the caller handle the error - grpcOvrlyLogger.Errorf("[ADD-WIRE:REMOTE-END] Remote end could not get interface index for %s. error:%v", hostEndVeth.LinkName, err) - return nil, err - } - grpcOvrlyLogger.Infof("[ADD-WIRE:REMOTE-END] Trigger from %s:%d : Successfully created remote pod to node vEth pair %s@%s <--> %s(%d).", - wireDef.PeerNodeIp, wireDef.WireIfIdOnPeerNode, inIfNm, wireDef.LocalPodName, outIfNm, locIface.Index) - aWire := CreateGWire(locIface.Index, hostEndVeth.LinkName, stopC, wireDef) - /* Utilizing google gopacket for polling for packets from the node. This seems to be the - simplest way to get all packets. - As an alternative to google gopacket(pcap), a socket based implementation is possible. - Not sure if socket based implementation can bring any advantage or not. - - Near term will replace pcap by socket. - */ - wrHandle, err := pcap.OpenLive(hostEndVeth.LinkName, 65365, true, pcap.BlockForever) - if err != nil { - // let the caller handle the error - grpcOvrlyLogger.Errorf("[ADD-WIRE:REMOTE-END] At remote end could not open interface (%d) for sed/recv packets for containers. error:%v", locIface.Index, err) - return nil, err - } + aWire := CreateGWire(int(wireID), wireDef.IntfNameInPod, stopC, wireDef) // Add the created wire in the in memory wire-map and k8S data store - AddWireInMemNDataStore(aWire, wrHandle) + AddWireInMemNDataStore(aWire, tapFile) return aWire, nil } diff --git a/third_party/meshnet/daemon/grpcwire/stream_manager.go b/third_party/meshnet/daemon/grpcwire/stream_manager.go new file mode 100644 index 00000000..6dcf0c9d --- /dev/null +++ b/third_party/meshnet/daemon/grpcwire/stream_manager.go @@ -0,0 +1,189 @@ +package grpcwire + +import ( + "context" + "fmt" + "strings" + "sync" + "time" + + "google.golang.org/grpc" + "google.golang.org/grpc/credentials/insecure" + + mpb "github.com/openconfig/kne/third_party/meshnet/daemon/proto/meshnet/v1beta1" + "github.com/openconfig/kne/third_party/meshnet/utils/wireutil" +) + +type nodeStreamKey struct { + topoNs string + peerIP string +} + +// NodeStream represents a single multiplexed gRPC packet stream shared by all TAP interfaces bound to a peer node within a topology. +type NodeStream struct { + key nodeStreamKey + pktChan chan *mpb.Packet + stopChan chan struct{} + refCount int + mu sync.Mutex + conn *grpc.ClientConn +} + +type nodeStreamManager struct { + mu sync.Mutex + streams map[nodeStreamKey]*NodeStream +} + +var streamMgr = &nodeStreamManager{ + streams: make(map[nodeStreamKey]*NodeStream), +} + +// GetOrCreateStream returns (or creates) the multiplexed NodeStream for the given topology namespace and peer IP. +// Increments refCount. +func (m *nodeStreamManager) GetOrCreateStream(topoNs string, peerIP string) *NodeStream { + if topoNs == "" { + topoNs = "default" + } + + m.mu.Lock() + defer m.mu.Unlock() + + key := nodeStreamKey{ + topoNs: topoNs, + peerIP: peerIP, + } + + if st, ok := m.streams[key]; ok { + st.refCount++ + return st + } + + st := &NodeStream{ + key: key, + pktChan: make(chan *mpb.Packet, 10000), + stopChan: make(chan struct{}), + refCount: 1, + } + m.streams[key] = st + + go st.run() + return st +} + +// ReleaseStream decrements refCount and stops the NodeStream if refCount reaches 0. +func (m *nodeStreamManager) ReleaseStream(topoNs string, peerIP string) { + if topoNs == "" { + topoNs = "default" + } + + m.mu.Lock() + defer m.mu.Unlock() + + key := nodeStreamKey{ + topoNs: topoNs, + peerIP: peerIP, + } + + st, ok := m.streams[key] + if !ok { + return + } + + st.refCount-- + if st.refCount <= 0 { + close(st.stopChan) + delete(m.streams, key) + } +} + +// Send enqueues a packet payload to be transmitted over the multiplexed gRPC stream. +func (s *NodeStream) Send(pkt *mpb.Packet) bool { + select { + case s.pktChan <- pkt: + return true + default: + // Queue full; drop packet under extreme overload + return false + } +} + +func (s *NodeStream) run() { + url := strings.TrimSpace(fmt.Sprintf("%s:%d", s.key.peerIP, wireutil.GRPCDefaultPort)) + dialOpts := []grpc.DialOption{ + grpc.WithTransportCredentials(insecure.NewCredentials()), + grpc.WithInitialWindowSize(4 * 1024 * 1024), // 4MB stream window + grpc.WithInitialConnWindowSize(16 * 1024 * 1024), // 16MB connection window + grpc.WithDefaultCallOptions( + grpc.MaxCallRecvMsgSize(64*1024*1024), + grpc.MaxCallSendMsgSize(64*1024*1024), + ), + } + + grpcOvrlyLogger.Infof("[STREAM-MGR] Starting multiplexed stream to peer %s in topo %s", s.key.peerIP, s.key.topoNs) + + for { + select { + case <-s.stopChan: + grpcOvrlyLogger.Infof("[STREAM-MGR] Stopping multiplexed stream to peer %s in topo %s", s.key.peerIP, s.key.topoNs) + return + default: + } + + conn, err := grpc.Dial(url, dialOpts...) + if err != nil { + grpcOvrlyLogger.Errorf("[STREAM-MGR] Failed to dial %s for topo %s: %v, retrying...", url, s.key.topoNs, err) + select { + case <-s.stopChan: + return + case <-time.After(2 * time.Second): + continue + } + } + + s.mu.Lock() + s.conn = conn + s.mu.Unlock() + + ctx, cancel := context.WithCancel(context.Background()) + client := mpb.NewWireProtocolClient(conn) + stream, err := client.SendToStream(ctx) + if err != nil { + cancel() + conn.Close() + grpcOvrlyLogger.Errorf("[STREAM-MGR] Failed to open SendToStream to %s for topo %s: %v, retrying...", url, s.key.topoNs, err) + select { + case <-s.stopChan: + return + case <-time.After(2 * time.Second): + continue + } + } + + grpcOvrlyLogger.Infof("[STREAM-MGR] Successfully connected multiplexed SendToStream to %s for topo %s", url, s.key.topoNs) + + s.drainAndSend(ctx, cancel, conn, stream) + } +} + +func (s *NodeStream) drainAndSend(ctx context.Context, cancel context.CancelFunc, conn *grpc.ClientConn, stream mpb.WireProtocol_SendToStreamClient) { + defer func() { + cancel() + _ = stream.CloseSend() + _ = conn.Close() + }() + + for { + select { + case <-s.stopChan: + return + case pkt, ok := <-s.pktChan: + if !ok { + return + } + if err := stream.Send(pkt); err != nil { + grpcOvrlyLogger.Debugf("[STREAM-MGR] Stream send error to %s for topo %s: %v", s.key.peerIP, s.key.topoNs, err) + return // break out to retry/reconnect loop in run() + } + } + } +} diff --git a/third_party/meshnet/daemon/main.go b/third_party/meshnet/daemon/main.go index 7ea4e7ac..8ee6dfcd 100644 --- a/third_party/meshnet/daemon/main.go +++ b/third_party/meshnet/daemon/main.go @@ -22,6 +22,8 @@ func main() { } defer cni.Cleanup() + wireutil.TuneSystem() + isDebug := flag.Bool("d", false, "enable degugging") grpcPort, err := strconv.Atoi(os.Getenv("GRPC_PORT")) if err != nil || grpcPort == 0 { diff --git a/third_party/meshnet/daemon/meshnet/controller.go b/third_party/meshnet/daemon/meshnet/controller.go index 1c6879b7..7da32cac 100644 --- a/third_party/meshnet/daemon/meshnet/controller.go +++ b/third_party/meshnet/daemon/meshnet/controller.go @@ -3,16 +3,13 @@ package meshnet import ( "context" "fmt" - "net" "strings" "time" - "github.com/containernetworking/plugins/pkg/ns" - mpb "github.com/openconfig/kne/third_party/meshnet/daemon/proto/meshnet/v1beta1" "github.com/openconfig/kne/third_party/meshnet/daemon/grpcwire" + mpb "github.com/openconfig/kne/third_party/meshnet/daemon/proto/meshnet/v1beta1" "github.com/openconfig/kne/third_party/meshnet/daemon/vxlan" "github.com/openconfig/kne/third_party/meshnet/utils/wireutil" - koko "github.com/redhat-nfvpe/koko/api" "github.com/vishvananda/netlink" "google.golang.org/grpc" "google.golang.org/grpc/credentials/insecure" @@ -129,6 +126,13 @@ func (m *Meshnet) reconcilePodLinksInternal(ctx context.Context, topo *unstructu } peerCache := make(map[string]*unstructured.Unstructured) + type grpcPeerBatch struct { + peerIP string + links []wireutil.PodLinkConfig + wireDefs []*mpb.WireDef + } + grpcBatches := make(map[string]*grpcPeerBatch) + sameNodeLinks := make([]wireutil.PodLinkConfig, 0, len(links)) for _, link := range links { peerTopo, ok := peerCache[link.PeerPodName] @@ -171,87 +175,35 @@ func (m *Meshnet) reconcilePodLinksInternal(ctx context.Context, topo *unstructu mnetdLogger.Infof("ReconcilePodLinks: initiating gRPC wire for pod %s <-> peer %s (UID %d)", topo.GetName(), link.PeerPodName, link.LinkUID) - // 1. Generate local host interface name - outIfNm, err := grpcwire.GenNodeIfaceName(topo.GetName(), link.LocalIntf) - if err != nil { - mnetdLogger.Errorf("ReconcilePodLinks: failed to generate node interface name: %v", err) - return err - } - - currNs, err := ns.GetCurrentNS() - if err != nil { - mnetdLogger.Errorf("ReconcilePodLinks: failed to get current host ns: %v", err) - return err - } - - // 2. Create VEth pair (pod namespace <-> host namespace) - inContainerVeth := koko.VEth{ - NsName: netNS, - LinkName: link.LocalIntf, - } - if link.LocalIP != "" { - ipAddr, ipSubnet, err := net.ParseCIDR(link.LocalIP) - if err != nil { - mnetdLogger.Errorf("ReconcilePodLinks: failed to parse local IP CIDR %s: %v", link.LocalIP, err) - return err - } - inContainerVeth.IPAddr = []net.IPNet{{ - IP: ipAddr, - Mask: ipSubnet.Mask, - }} - } - - hostEndVeth := koko.VEth{ - NsName: currNs.Path(), - LinkName: outIfNm, - } - - // Make Veth - if err := koko.MakeVeth(inContainerVeth, hostEndVeth); err != nil { - mnetdLogger.Errorf("ReconcilePodLinks: failed to create local vEth pair (in:%s, out:%s) for gRPC wire: %v", - inContainerVeth.LinkName, hostEndVeth.LinkName, err) - return err - } - - // Disable TX checksum - if err := wireutil.SetTxChecksumOff(inContainerVeth.LinkName, inContainerVeth.NsName); err != nil { - mnetdLogger.Errorf("ReconcilePodLinks: failed to disable Tx checksum on %s: %v", inContainerVeth.LinkName, err) - } - - // 3. Register local end in meshnet daemon (acts as CreateGRPCWireLocal) + // 1. Register local end in meshnet daemon (creates/attaches TAP interface in container netns) wireDefLocal := &mpb.WireDef{ - LocalPodNetNs: netNS, - LinkUid: link.LinkUID, - TopoNs: link.KubeNs, - WireIfNameOnLocalNode: outIfNm, - LocalPodName: topo.GetName(), - IntfNameInPod: link.LocalIntf, - LocalPodIp: link.LocalIP, - PeerNodeIp: peerSrcIP, + LocalPodNetNs: netNS, + LinkUid: link.LinkUID, + TopoNs: link.KubeNs, + LocalPodName: topo.GetName(), + IntfNameInPod: link.LocalIntf, + LocalPodIp: link.LocalIP, + PeerNodeIp: peerSrcIP, } if _, err := grpcwire.CreateGRPCWireLocal(ctx, wireDefLocal); err != nil { mnetdLogger.Errorf("ReconcilePodLinks: failed to register local GRPC wire: %v", err) return err } - // 4. Dial remote daemon and trigger remote end creation - url := fmt.Sprintf("%s:%d", peerSrcIP, wireutil.GRPCDefaultPort) - url = strings.TrimSpace(url) - remoteConn, err := grpc.Dial(url, grpc.WithTransportCredentials(insecure.NewCredentials())) - if err != nil { - mnetdLogger.Errorf("ReconcilePodLinks: failed to dial remote node %s: %v", url, err) - return err + locWire, ok := grpcwire.GetWireByUID(netNS, int(link.LinkUID)) + if !ok || locWire == nil { + mnetdLogger.Errorf("ReconcilePodLinks: failed to get local wire for link UID %d", link.LinkUID) + return fmt.Errorf("local wire not found for link UID %d", link.LinkUID) } - locInf, err := net.InterfaceByName(outIfNm) - if err != nil { - remoteConn.Close() - mnetdLogger.Errorf("ReconcilePodLinks: failed to get local interface by name %s: %v", outIfNm, err) - return err + pBatch := grpcBatches[peerSrcIP] + if pBatch == nil { + pBatch = &grpcPeerBatch{peerIP: peerSrcIP} + grpcBatches[peerSrcIP] = pBatch } - - wireDefRemote := &mpb.WireDef{ - WireIfIdOnPeerNode: int64(locInf.Index), + pBatch.links = append(pBatch.links, link) + pBatch.wireDefs = append(pBatch.wireDefs, &mpb.WireDef{ + WireIfIdOnPeerNode: locWire.LocalNodeIfaceID, PeerNodeIp: srcIP, IntfNameInPod: link.PeerIntf, LocalPodNetNs: peerNetNS, @@ -259,20 +211,7 @@ func (m *Meshnet) reconcilePodLinksInternal(ctx context.Context, topo *unstructu LinkUid: link.LinkUID, TopoNs: link.KubeNs, LocalPodIp: link.PeerIP, - } - - remoteClient := mpb.NewRemoteClient(remoteConn) - mnetdLogger.Infof("ReconcilePodLinks: calling remote node AddGRPCWireRemote (%s) for link UID %d", url, link.LinkUID) - creatResp, err := remoteClient.AddGRPCWireRemote(ctx, wireDefRemote) - if err != nil || !creatResp.Response { - remoteConn.Close() - mnetdLogger.Errorf("ReconcilePodLinks: remote AddGRPCWireRemote failed: %v", err) - return fmt.Errorf("remote AddGRPCWireRemote failed: %v", err) - } - remoteConn.Close() - - // 5. Update local end with the peer's host interface ID returned by Node 2 - grpcwire.UpdateWireByUID(netNS, int(link.LinkUID), creatResp.PeerIntfId, make(chan struct{})) + }) } else { remotePod := &mpb.RemotePod{ NetNs: netNS, @@ -292,6 +231,49 @@ func (m *Meshnet) reconcilePodLinksInternal(ctx context.Context, topo *unstructu } } + // 2. Process gRPC peer batches in chunks (default 50 items per RPC) to allow pipelined processing + batchSize := wireutil.GetEnvInt("WIRE_BATCH_SIZE", 50) + + for peerIP, pBatch := range grpcBatches { + url := fmt.Sprintf("%s:%d", peerIP, wireutil.GRPCDefaultPort) + url = strings.TrimSpace(url) + remoteConn, err := grpc.Dial(url, grpc.WithTransportCredentials(insecure.NewCredentials())) + if err != nil { + mnetdLogger.Errorf("ReconcilePodLinks: failed to dial remote node %s: %v", url, err) + return err + } + + remoteClient := mpb.NewRemoteClient(remoteConn) + total := len(pBatch.wireDefs) + + for i := 0; i < total; i += batchSize { + end := i + batchSize + if end > total { + end = total + } + + chunkWireDefs := pBatch.wireDefs[i:end] + chunkLinks := pBatch.links[i:end] + + mnetdLogger.Infof("ReconcilePodLinks: calling AddGRPCWiresRemoteBatch on %s for batch [%d:%d] of %d links", url, i, end, total) + batchResp, err := remoteClient.AddGRPCWiresRemoteBatch(ctx, &mpb.WireDefBatch{Items: chunkWireDefs}) + if err != nil { + remoteConn.Close() + mnetdLogger.Errorf("ReconcilePodLinks: AddGRPCWiresRemoteBatch failed to %s for batch [%d:%d]: %v", url, i, end, err) + return fmt.Errorf("AddGRPCWiresRemoteBatch failed: %v", err) + } + + for j, res := range batchResp.Items { + if res != nil && res.Response { + l := chunkLinks[j] + grpcwire.UpdateWireByUID(netNS, int(l.LinkUID), res.PeerIntfId, make(chan struct{})) + } + } + } + + remoteConn.Close() + } + if len(sameNodeLinks) > 0 { mnetdLogger.Infof("ReconcilePodLinks: configuring %d active same-node links for pod %s (%s)", len(sameNodeLinks), topo.GetName(), netNS) if err := wireutil.ConfigurePodLinks(netNS, sameNodeLinks); err != nil { diff --git a/third_party/meshnet/daemon/meshnet/handler.go b/third_party/meshnet/daemon/meshnet/handler.go index ac18e1d1..01801c83 100644 --- a/third_party/meshnet/daemon/meshnet/handler.go +++ b/third_party/meshnet/daemon/meshnet/handler.go @@ -3,6 +3,7 @@ package meshnet import ( "context" "fmt" + "io" "os" "github.com/openconfig/kne/third_party/meshnet/api/types/v1beta1" @@ -380,6 +381,14 @@ func (m *Meshnet) AddGRPCWireLocal(ctx context.Context, wireDef *mpb.WireDef) (* // ------------------------------------------------------------------------------------------------------ func (m *Meshnet) SendToOnce(ctx context.Context, pkt *mpb.Packet) (*mpb.BoolResponse, error) { + if pkt.RemotIntfId <= 0 { + log.WithFields(log.Fields{ + "daemon": "meshnetd", + "overlay": "gRPC", + }).Debugf("SendToOnce: received packet for uninitialized wire id %d, peer not ready yet", pkt.RemotIntfId) + return &mpb.BoolResponse{Response: false}, nil + } + wrHandle, err := grpcwire.GetHostIntfHndl(pkt.RemotIntfId) if err != nil { log.WithFields(log.Fields{ @@ -394,7 +403,7 @@ func (m *Meshnet) SendToOnce(ctx context.Context, pkt *mpb.Packet) (*mpb.BoolRes // log.Printf("Daemon(SendToOnce): Received [pkt: %s, bytes: %d, for local interface id: %d]. Sending it to local container", pktType, len(pkt.Frame), pkt.RemotIntfId) // log.Printf("Daemon(SendToOnce): Received [bytes: %d, for local interface id: %d]. Sending it to local container", len(pkt.Frame), pkt.RemotIntfId) - err = wrHandle.WritePacketData(pkt.Frame) + _, err = wrHandle.Write(pkt.Frame) if err != nil { log.WithFields(log.Fields{ "daemon": "meshnetd", @@ -406,6 +415,39 @@ func (m *Meshnet) SendToOnce(ctx context.Context, pkt *mpb.Packet) (*mpb.BoolRes return &mpb.BoolResponse{Response: true}, nil } +// ------------------------------------------------------------------------------------------------------ +func (m *Meshnet) SendToStream(stream mpb.WireProtocol_SendToStreamServer) error { + for { + pkt, err := stream.Recv() + if err == io.EOF { + return stream.SendAndClose(&mpb.BoolResponse{Response: true}) + } + if err != nil { + return err + } + + if pkt.RemotIntfId <= 0 { + continue + } + + wrHandle, err := grpcwire.GetHostIntfHndl(pkt.RemotIntfId) + if err != nil { + log.WithFields(log.Fields{ + "daemon": "meshnetd", + "overlay": "gRPC", + }).Debugf("SendToStream (wire id - %v): Could not find local handle. err:%v", pkt.RemotIntfId, err) + continue + } + + if _, err := wrHandle.Write(pkt.Frame); err != nil { + log.WithFields(log.Fields{ + "daemon": "meshnetd", + "overlay": "gRPC", + }).Errorf("SendToStream (wire id - %v): Could not write packet(%d bytes) to local interface. err:%v", pkt.RemotIntfId, len(pkt.Frame), err) + } + } +} + // --------------------------------------------------------------------------------------------------------------- func (m *Meshnet) AddGRPCWireRemote(ctx context.Context, wireDef *mpb.WireDef) (*mpb.WireCreateResponse, error) { stopC := make(chan struct{}) @@ -426,6 +468,31 @@ func (m *Meshnet) AddGRPCWireRemote(ctx context.Context, wireDef *mpb.WireDef) ( return &mpb.WireCreateResponse{Response: false, PeerIntfId: wireDef.WireIfIdOnPeerNode}, err } +// AddGRPCWiresRemoteBatch handles batch creation of remote gRPC wires in a single RPC call. +func (m *Meshnet) AddGRPCWiresRemoteBatch(ctx context.Context, req *mpb.WireDefBatch) (*mpb.WireCreateResponseBatch, error) { + if req == nil { + return &mpb.WireCreateResponseBatch{}, nil + } + resp := &mpb.WireCreateResponseBatch{ + Items: make([]*mpb.WireCreateResponse, len(req.Items)), + } + for i, wireDef := range req.Items { + stopC := make(chan struct{}) + wire, err := grpcwire.CreateUpdateGRPCWireRemoteTriggered(wireDef, stopC) + if err != nil { + log.WithFields(log.Fields{ + "daemon": "meshnetd", + "overlay": "gRPC", + }).Errorf("[ADD-WIRE:REMOTE-END-BATCH] Error creating wire %s@%s: %v", wireDef.LocalPodName, wireDef.IntfNameInPod, err) + resp.Items[i] = &mpb.WireCreateResponse{Response: false} + continue + } + go grpcwire.RecvFrmLocalPodThread(wire, wire.LocalNodeIfaceName) + resp.Items[i] = &mpb.WireCreateResponse{Response: true, PeerIntfId: wire.LocalNodeIfaceID} + } + return resp, nil +} + // --------------------------------------------------------------------------------------------------------------- func (m *Meshnet) GRPCWireDownRemote(ctx context.Context, wireDef *mpb.WireDef) (*mpb.WireDownResponse, error) { err := grpcwire.GRPCWireDownRemoteTriggered(wireDef) @@ -447,10 +514,10 @@ func (m *Meshnet) GRPCWireDownRemote(ctx context.Context, wireDef *mpb.WireDef) // GRPCWireExists will return the wire if it exists. func (m *Meshnet) GRPCWireExists(ctx context.Context, wireDef *mpb.WireDef) (*mpb.WireCreateResponse, error) { wire, ok := grpcwire.GetWireByUID(wireDef.LocalPodNetNs, int(wireDef.LinkUid)) - if !ok || wire == nil { + if !ok || wire == nil || !wire.IsReady || wire.WireIfaceIDOnPeerNode <= 0 { return &mpb.WireCreateResponse{Response: false, PeerIntfId: wireDef.WireIfIdOnPeerNode}, nil } - return &mpb.WireCreateResponse{Response: ok, PeerIntfId: wire.WireIfaceIDOnPeerNode}, nil + return &mpb.WireCreateResponse{Response: true, PeerIntfId: wire.WireIfaceIDOnPeerNode}, nil } // --------------------------------------------------------------------------------------------------------------- diff --git a/third_party/meshnet/daemon/meshnet/meshnet.go b/third_party/meshnet/daemon/meshnet/meshnet.go index 341764a3..065ebf94 100644 --- a/third_party/meshnet/daemon/meshnet/meshnet.go +++ b/third_party/meshnet/daemon/meshnet/meshnet.go @@ -1,3 +1,5 @@ +// Package meshnet implements the meshnet daemon controller loop, K8s topology resource watching, +// and gRPC wire/vxLAN link reconciliation. package meshnet import ( @@ -27,22 +29,24 @@ import ( mpb "github.com/openconfig/kne/third_party/meshnet/daemon/proto/meshnet/v1beta1" ) +// Config defines configuration options for initializing the Meshnet daemon server. type Config struct { Port int GRPCOpts []grpc.ServerOption } +// Meshnet represents the main daemon service instance handling Kubernetes topology reconciliation and gRPC wire protocol RPCs. type Meshnet struct { mpb.UnimplementedLocalServer mpb.UnimplementedRemoteServer mpb.UnimplementedWireProtocolServer - config Config - kClient kubernetes.Interface - tClient topologyclientv1.Interface - GWireDynClient *dynamic.DynamicClient - rCfg *rest.Config - s *grpc.Server - lis net.Listener + config Config + kClient kubernetes.Interface + tClient topologyclientv1.Interface + GWireDynClient *dynamic.DynamicClient + rCfg *rest.Config + s *grpc.Server + lis net.Listener nodeIP string dirtyChan chan struct{} interNodeLinkType string @@ -50,6 +54,7 @@ type Meshnet struct { var mnetdLogger *log.Entry = nil +// InitLogger initializes the logrus logger for the meshnet daemon. func InitLogger() { mnetdLogger = log.WithFields(log.Fields{"daemon": "meshnetd"}) } @@ -72,6 +77,7 @@ func restConfig() (*rest.Config, error) { return rCfg, nil } +// New creates and initializes a new Meshnet daemon instance with gRPC server options and K8s clientsets. func New(cfg Config) (*Meshnet, error) { rCfg, err := restConfig() if err != nil { @@ -97,25 +103,33 @@ func New(cfg Config) (*Meshnet, error) { // If the link type is GRPC then set the GRPC logging level to LevelNone // Otherwise there will be GRPC log for every packet sent as for link type GRPC, GRPC is also the data-plane. This is too // much of log that does not help in debugging and K8S does log rotation very frequently. + defaultOpts := []grpc.ServerOption{ + grpc.InitialWindowSize(4 * 1024 * 1024), // 4MB stream window + grpc.InitialConnWindowSize(16 * 1024 * 1024), // 16MB connection window + grpc.MaxRecvMsgSize(64 * 1024 * 1024), + grpc.MaxSendMsgSize(64 * 1024 * 1024), + } + allOpts := append(defaultOpts, cfg.GRPCOpts...) + var svr *grpc.Server lnkTyp := os.Getenv("INTER_NODE_LINK_TYPE") if lnkTyp == wireutil.INTER_NODE_LINK_GRPC { - svr = grpc.NewServer(cfg.GRPCOpts...) + svr = grpc.NewServer(allOpts...) } else { - svr = newServerWithLogging(cfg.GRPCOpts...) + svr = newServerWithLogging(allOpts...) } m := &Meshnet{ - config: cfg, - rCfg: rCfg, - kClient: kClient, - tClient: tClient, - GWireDynClient: gwireDynClient, - lis: lis, - s: svr, + config: cfg, + rCfg: rCfg, + kClient: kClient, + tClient: tClient, + GWireDynClient: gwireDynClient, + lis: lis, + s: svr, nodeIP: os.Getenv("HOST_IP"), dirtyChan: make(chan struct{}, 1), - interNodeLinkType: lnkTyp, + interNodeLinkType: lnkTyp, } mpb.RegisterLocalServer(m.s, m) mpb.RegisterRemoteServer(m.s, m) @@ -134,11 +148,13 @@ func New(cfg Config) (*Meshnet, error) { return m, nil } +// Serve starts the gRPC server listening on the configured port. func (m *Meshnet) Serve() error { mnetdLogger.Infof("GRPC server has started on port: %d", m.config.Port) return m.s.Serve(m.lis) } +// Stop gracefully stops the gRPC server instance. func (m *Meshnet) Stop() { m.s.Stop() } diff --git a/third_party/meshnet/daemon/proto/meshnet/v1beta1/meshnet.pb.go b/third_party/meshnet/daemon/proto/meshnet/v1beta1/meshnet.pb.go index 973f989f..ee86bcdd 100644 --- a/third_party/meshnet/daemon/proto/meshnet/v1beta1/meshnet.pb.go +++ b/third_party/meshnet/daemon/proto/meshnet/v1beta1/meshnet.pb.go @@ -643,6 +643,94 @@ func (x *WireCreateResponse) GetPeerIntfId() int64 { return 0 } +type WireDefBatch struct { + state protoimpl.MessageState `protogen:"open.v1"` + Items []*WireDef `protobuf:"bytes,1,rep,name=items,proto3" json:"items,omitempty"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *WireDefBatch) Reset() { + *x = WireDefBatch{} + mi := &file_daemon_proto_meshnet_v1beta1_meshnet_proto_msgTypes[8] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *WireDefBatch) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*WireDefBatch) ProtoMessage() {} + +func (x *WireDefBatch) ProtoReflect() protoreflect.Message { + mi := &file_daemon_proto_meshnet_v1beta1_meshnet_proto_msgTypes[8] + if x != nil { + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + if ms.LoadMessageInfo() == nil { + ms.StoreMessageInfo(mi) + } + return ms + } + return mi.MessageOf(x) +} + +// Deprecated: Use WireDefBatch.ProtoReflect.Descriptor instead. +func (*WireDefBatch) Descriptor() ([]byte, []int) { + return file_daemon_proto_meshnet_v1beta1_meshnet_proto_rawDescGZIP(), []int{8} +} + +func (x *WireDefBatch) GetItems() []*WireDef { + if x != nil { + return x.Items + } + return nil +} + +type WireCreateResponseBatch struct { + state protoimpl.MessageState `protogen:"open.v1"` + Items []*WireCreateResponse `protobuf:"bytes,1,rep,name=items,proto3" json:"items,omitempty"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *WireCreateResponseBatch) Reset() { + *x = WireCreateResponseBatch{} + mi := &file_daemon_proto_meshnet_v1beta1_meshnet_proto_msgTypes[9] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *WireCreateResponseBatch) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*WireCreateResponseBatch) ProtoMessage() {} + +func (x *WireCreateResponseBatch) ProtoReflect() protoreflect.Message { + mi := &file_daemon_proto_meshnet_v1beta1_meshnet_proto_msgTypes[9] + if x != nil { + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + if ms.LoadMessageInfo() == nil { + ms.StoreMessageInfo(mi) + } + return ms + } + return mi.MessageOf(x) +} + +// Deprecated: Use WireCreateResponseBatch.ProtoReflect.Descriptor instead. +func (*WireCreateResponseBatch) Descriptor() ([]byte, []int) { + return file_daemon_proto_meshnet_v1beta1_meshnet_proto_rawDescGZIP(), []int{9} +} + +func (x *WireCreateResponseBatch) GetItems() []*WireCreateResponse { + if x != nil { + return x.Items + } + return nil +} + type WireDownResponse struct { state protoimpl.MessageState `protogen:"open.v1"` Response bool `protobuf:"varint,1,opt,name=response,proto3" json:"response,omitempty"` @@ -652,7 +740,7 @@ type WireDownResponse struct { func (x *WireDownResponse) Reset() { *x = WireDownResponse{} - mi := &file_daemon_proto_meshnet_v1beta1_meshnet_proto_msgTypes[8] + mi := &file_daemon_proto_meshnet_v1beta1_meshnet_proto_msgTypes[10] ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) ms.StoreMessageInfo(mi) } @@ -664,7 +752,7 @@ func (x *WireDownResponse) String() string { func (*WireDownResponse) ProtoMessage() {} func (x *WireDownResponse) ProtoReflect() protoreflect.Message { - mi := &file_daemon_proto_meshnet_v1beta1_meshnet_proto_msgTypes[8] + mi := &file_daemon_proto_meshnet_v1beta1_meshnet_proto_msgTypes[10] if x != nil { ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) if ms.LoadMessageInfo() == nil { @@ -677,7 +765,7 @@ func (x *WireDownResponse) ProtoReflect() protoreflect.Message { // Deprecated: Use WireDownResponse.ProtoReflect.Descriptor instead. func (*WireDownResponse) Descriptor() ([]byte, []int) { - return file_daemon_proto_meshnet_v1beta1_meshnet_proto_rawDescGZIP(), []int{8} + return file_daemon_proto_meshnet_v1beta1_meshnet_proto_rawDescGZIP(), []int{10} } func (x *WireDownResponse) GetResponse() bool { @@ -698,7 +786,7 @@ type Packet struct { func (x *Packet) Reset() { *x = Packet{} - mi := &file_daemon_proto_meshnet_v1beta1_meshnet_proto_msgTypes[9] + mi := &file_daemon_proto_meshnet_v1beta1_meshnet_proto_msgTypes[11] ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) ms.StoreMessageInfo(mi) } @@ -710,7 +798,7 @@ func (x *Packet) String() string { func (*Packet) ProtoMessage() {} func (x *Packet) ProtoReflect() protoreflect.Message { - mi := &file_daemon_proto_meshnet_v1beta1_meshnet_proto_msgTypes[9] + mi := &file_daemon_proto_meshnet_v1beta1_meshnet_proto_msgTypes[11] if x != nil { ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) if ms.LoadMessageInfo() == nil { @@ -723,7 +811,7 @@ func (x *Packet) ProtoReflect() protoreflect.Message { // Deprecated: Use Packet.ProtoReflect.Descriptor instead. func (*Packet) Descriptor() ([]byte, []int) { - return file_daemon_proto_meshnet_v1beta1_meshnet_proto_rawDescGZIP(), []int{9} + return file_daemon_proto_meshnet_v1beta1_meshnet_proto_rawDescGZIP(), []int{11} } func (x *Packet) GetRemotIntfId() int64 { @@ -750,7 +838,7 @@ type GenerateNodeInterfaceNameRequest struct { func (x *GenerateNodeInterfaceNameRequest) Reset() { *x = GenerateNodeInterfaceNameRequest{} - mi := &file_daemon_proto_meshnet_v1beta1_meshnet_proto_msgTypes[10] + mi := &file_daemon_proto_meshnet_v1beta1_meshnet_proto_msgTypes[12] ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) ms.StoreMessageInfo(mi) } @@ -762,7 +850,7 @@ func (x *GenerateNodeInterfaceNameRequest) String() string { func (*GenerateNodeInterfaceNameRequest) ProtoMessage() {} func (x *GenerateNodeInterfaceNameRequest) ProtoReflect() protoreflect.Message { - mi := &file_daemon_proto_meshnet_v1beta1_meshnet_proto_msgTypes[10] + mi := &file_daemon_proto_meshnet_v1beta1_meshnet_proto_msgTypes[12] if x != nil { ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) if ms.LoadMessageInfo() == nil { @@ -775,7 +863,7 @@ func (x *GenerateNodeInterfaceNameRequest) ProtoReflect() protoreflect.Message { // Deprecated: Use GenerateNodeInterfaceNameRequest.ProtoReflect.Descriptor instead. func (*GenerateNodeInterfaceNameRequest) Descriptor() ([]byte, []int) { - return file_daemon_proto_meshnet_v1beta1_meshnet_proto_rawDescGZIP(), []int{10} + return file_daemon_proto_meshnet_v1beta1_meshnet_proto_rawDescGZIP(), []int{12} } func (x *GenerateNodeInterfaceNameRequest) GetPodIntfName() string { @@ -802,7 +890,7 @@ type GenerateNodeInterfaceNameResponse struct { func (x *GenerateNodeInterfaceNameResponse) Reset() { *x = GenerateNodeInterfaceNameResponse{} - mi := &file_daemon_proto_meshnet_v1beta1_meshnet_proto_msgTypes[11] + mi := &file_daemon_proto_meshnet_v1beta1_meshnet_proto_msgTypes[13] ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) ms.StoreMessageInfo(mi) } @@ -814,7 +902,7 @@ func (x *GenerateNodeInterfaceNameResponse) String() string { func (*GenerateNodeInterfaceNameResponse) ProtoMessage() {} func (x *GenerateNodeInterfaceNameResponse) ProtoReflect() protoreflect.Message { - mi := &file_daemon_proto_meshnet_v1beta1_meshnet_proto_msgTypes[11] + mi := &file_daemon_proto_meshnet_v1beta1_meshnet_proto_msgTypes[13] if x != nil { ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) if ms.LoadMessageInfo() == nil { @@ -827,7 +915,7 @@ func (x *GenerateNodeInterfaceNameResponse) ProtoReflect() protoreflect.Message // Deprecated: Use GenerateNodeInterfaceNameResponse.ProtoReflect.Descriptor instead. func (*GenerateNodeInterfaceNameResponse) Descriptor() ([]byte, []int) { - return file_daemon_proto_meshnet_v1beta1_meshnet_proto_rawDescGZIP(), []int{11} + return file_daemon_proto_meshnet_v1beta1_meshnet_proto_rawDescGZIP(), []int{13} } func (x *GenerateNodeInterfaceNameResponse) GetOk() bool { @@ -899,7 +987,11 @@ const file_daemon_proto_meshnet_v1beta1_meshnet_proto_rawDesc = "" + "\x12WireCreateResponse\x12\x1a\n" + "\bresponse\x18\x01 \x01(\bR\bresponse\x12 \n" + "\fpeer_intf_id\x18\x02 \x01(\x03R\n" + - "peerIntfId\".\n" + + "peerIntfId\">\n" + + "\fWireDefBatch\x12.\n" + + "\x05items\x18\x01 \x03(\v2\x18.meshnet.v1beta1.WireDefR\x05items\"T\n" + + "\x17WireCreateResponseBatch\x129\n" + + "\x05items\x18\x01 \x03(\v2#.meshnet.v1beta1.WireCreateResponseR\x05items\".\n" + "\x10WireDownResponse\x12\x1a\n" + "\bresponse\x18\x01 \x01(\bR\bresponse\"B\n" + "\x06Packet\x12\"\n" + @@ -920,10 +1012,11 @@ const file_daemon_proto_meshnet_v1beta1_meshnet_proto_rawDesc = "" + "\x0eGRPCWireExists\x12\x18.meshnet.v1beta1.WireDef\x1a#.meshnet.v1beta1.WireCreateResponse\x12K\n" + "\x10AddGRPCWireLocal\x12\x18.meshnet.v1beta1.WireDef\x1a\x1d.meshnet.v1beta1.BoolResponse\x12F\n" + "\vRemGRPCWire\x12\x18.meshnet.v1beta1.WireDef\x1a\x1d.meshnet.v1beta1.BoolResponse\x12\x82\x01\n" + - "\x19GenerateNodeInterfaceName\x121.meshnet.v1beta1.GenerateNodeInterfaceNameRequest\x1a2.meshnet.v1beta1.GenerateNodeInterfaceNameResponse2\xf4\x01\n" + + "\x19GenerateNodeInterfaceName\x121.meshnet.v1beta1.GenerateNodeInterfaceNameRequest\x1a2.meshnet.v1beta1.GenerateNodeInterfaceNameResponse2\xd8\x02\n" + "\x06Remote\x12C\n" + "\x06Update\x12\x1a.meshnet.v1beta1.RemotePod\x1a\x1d.meshnet.v1beta1.BoolResponse\x12R\n" + - "\x11AddGRPCWireRemote\x12\x18.meshnet.v1beta1.WireDef\x1a#.meshnet.v1beta1.WireCreateResponse\x12Q\n" + + "\x11AddGRPCWireRemote\x12\x18.meshnet.v1beta1.WireDef\x1a#.meshnet.v1beta1.WireCreateResponse\x12b\n" + + "\x17AddGRPCWiresRemoteBatch\x12\x1d.meshnet.v1beta1.WireDefBatch\x1a(.meshnet.v1beta1.WireCreateResponseBatch\x12Q\n" + "\x12GRPCWireDownRemote\x12\x18.meshnet.v1beta1.WireDef\x1a!.meshnet.v1beta1.WireDownResponse2\x9e\x01\n" + "\fWireProtocol\x12D\n" + "\n" + @@ -942,7 +1035,7 @@ func file_daemon_proto_meshnet_v1beta1_meshnet_proto_rawDescGZIP() []byte { return file_daemon_proto_meshnet_v1beta1_meshnet_proto_rawDescData } -var file_daemon_proto_meshnet_v1beta1_meshnet_proto_msgTypes = make([]protoimpl.MessageInfo, 12) +var file_daemon_proto_meshnet_v1beta1_meshnet_proto_msgTypes = make([]protoimpl.MessageInfo, 14) var file_daemon_proto_meshnet_v1beta1_meshnet_proto_goTypes = []any{ (*Pod)(nil), // 0: meshnet.v1beta1.Pod (*Link)(nil), // 1: meshnet.v1beta1.Link @@ -952,46 +1045,52 @@ var file_daemon_proto_meshnet_v1beta1_meshnet_proto_goTypes = []any{ (*RemotePod)(nil), // 5: meshnet.v1beta1.RemotePod (*WireDef)(nil), // 6: meshnet.v1beta1.WireDef (*WireCreateResponse)(nil), // 7: meshnet.v1beta1.WireCreateResponse - (*WireDownResponse)(nil), // 8: meshnet.v1beta1.WireDownResponse - (*Packet)(nil), // 9: meshnet.v1beta1.Packet - (*GenerateNodeInterfaceNameRequest)(nil), // 10: meshnet.v1beta1.GenerateNodeInterfaceNameRequest - (*GenerateNodeInterfaceNameResponse)(nil), // 11: meshnet.v1beta1.GenerateNodeInterfaceNameResponse + (*WireDefBatch)(nil), // 8: meshnet.v1beta1.WireDefBatch + (*WireCreateResponseBatch)(nil), // 9: meshnet.v1beta1.WireCreateResponseBatch + (*WireDownResponse)(nil), // 10: meshnet.v1beta1.WireDownResponse + (*Packet)(nil), // 11: meshnet.v1beta1.Packet + (*GenerateNodeInterfaceNameRequest)(nil), // 12: meshnet.v1beta1.GenerateNodeInterfaceNameRequest + (*GenerateNodeInterfaceNameResponse)(nil), // 13: meshnet.v1beta1.GenerateNodeInterfaceNameResponse } var file_daemon_proto_meshnet_v1beta1_meshnet_proto_depIdxs = []int32{ 1, // 0: meshnet.v1beta1.Pod.links:type_name -> meshnet.v1beta1.Link - 2, // 1: meshnet.v1beta1.Local.Get:input_type -> meshnet.v1beta1.PodQuery - 0, // 2: meshnet.v1beta1.Local.SetAlive:input_type -> meshnet.v1beta1.Pod - 3, // 3: meshnet.v1beta1.Local.SkipReverse:input_type -> meshnet.v1beta1.SkipQuery - 3, // 4: meshnet.v1beta1.Local.Skip:input_type -> meshnet.v1beta1.SkipQuery - 3, // 5: meshnet.v1beta1.Local.IsSkipped:input_type -> meshnet.v1beta1.SkipQuery - 6, // 6: meshnet.v1beta1.Local.GRPCWireExists:input_type -> meshnet.v1beta1.WireDef - 6, // 7: meshnet.v1beta1.Local.AddGRPCWireLocal:input_type -> meshnet.v1beta1.WireDef - 6, // 8: meshnet.v1beta1.Local.RemGRPCWire:input_type -> meshnet.v1beta1.WireDef - 10, // 9: meshnet.v1beta1.Local.GenerateNodeInterfaceName:input_type -> meshnet.v1beta1.GenerateNodeInterfaceNameRequest - 5, // 10: meshnet.v1beta1.Remote.Update:input_type -> meshnet.v1beta1.RemotePod - 6, // 11: meshnet.v1beta1.Remote.AddGRPCWireRemote:input_type -> meshnet.v1beta1.WireDef - 6, // 12: meshnet.v1beta1.Remote.GRPCWireDownRemote:input_type -> meshnet.v1beta1.WireDef - 9, // 13: meshnet.v1beta1.WireProtocol.SendToOnce:input_type -> meshnet.v1beta1.Packet - 9, // 14: meshnet.v1beta1.WireProtocol.SendToStream:input_type -> meshnet.v1beta1.Packet - 0, // 15: meshnet.v1beta1.Local.Get:output_type -> meshnet.v1beta1.Pod - 4, // 16: meshnet.v1beta1.Local.SetAlive:output_type -> meshnet.v1beta1.BoolResponse - 4, // 17: meshnet.v1beta1.Local.SkipReverse:output_type -> meshnet.v1beta1.BoolResponse - 4, // 18: meshnet.v1beta1.Local.Skip:output_type -> meshnet.v1beta1.BoolResponse - 4, // 19: meshnet.v1beta1.Local.IsSkipped:output_type -> meshnet.v1beta1.BoolResponse - 7, // 20: meshnet.v1beta1.Local.GRPCWireExists:output_type -> meshnet.v1beta1.WireCreateResponse - 4, // 21: meshnet.v1beta1.Local.AddGRPCWireLocal:output_type -> meshnet.v1beta1.BoolResponse - 4, // 22: meshnet.v1beta1.Local.RemGRPCWire:output_type -> meshnet.v1beta1.BoolResponse - 11, // 23: meshnet.v1beta1.Local.GenerateNodeInterfaceName:output_type -> meshnet.v1beta1.GenerateNodeInterfaceNameResponse - 4, // 24: meshnet.v1beta1.Remote.Update:output_type -> meshnet.v1beta1.BoolResponse - 7, // 25: meshnet.v1beta1.Remote.AddGRPCWireRemote:output_type -> meshnet.v1beta1.WireCreateResponse - 8, // 26: meshnet.v1beta1.Remote.GRPCWireDownRemote:output_type -> meshnet.v1beta1.WireDownResponse - 4, // 27: meshnet.v1beta1.WireProtocol.SendToOnce:output_type -> meshnet.v1beta1.BoolResponse - 4, // 28: meshnet.v1beta1.WireProtocol.SendToStream:output_type -> meshnet.v1beta1.BoolResponse - 15, // [15:29] is the sub-list for method output_type - 1, // [1:15] is the sub-list for method input_type - 1, // [1:1] is the sub-list for extension type_name - 1, // [1:1] is the sub-list for extension extendee - 0, // [0:1] is the sub-list for field type_name + 6, // 1: meshnet.v1beta1.WireDefBatch.items:type_name -> meshnet.v1beta1.WireDef + 7, // 2: meshnet.v1beta1.WireCreateResponseBatch.items:type_name -> meshnet.v1beta1.WireCreateResponse + 2, // 3: meshnet.v1beta1.Local.Get:input_type -> meshnet.v1beta1.PodQuery + 0, // 4: meshnet.v1beta1.Local.SetAlive:input_type -> meshnet.v1beta1.Pod + 3, // 5: meshnet.v1beta1.Local.SkipReverse:input_type -> meshnet.v1beta1.SkipQuery + 3, // 6: meshnet.v1beta1.Local.Skip:input_type -> meshnet.v1beta1.SkipQuery + 3, // 7: meshnet.v1beta1.Local.IsSkipped:input_type -> meshnet.v1beta1.SkipQuery + 6, // 8: meshnet.v1beta1.Local.GRPCWireExists:input_type -> meshnet.v1beta1.WireDef + 6, // 9: meshnet.v1beta1.Local.AddGRPCWireLocal:input_type -> meshnet.v1beta1.WireDef + 6, // 10: meshnet.v1beta1.Local.RemGRPCWire:input_type -> meshnet.v1beta1.WireDef + 12, // 11: meshnet.v1beta1.Local.GenerateNodeInterfaceName:input_type -> meshnet.v1beta1.GenerateNodeInterfaceNameRequest + 5, // 12: meshnet.v1beta1.Remote.Update:input_type -> meshnet.v1beta1.RemotePod + 6, // 13: meshnet.v1beta1.Remote.AddGRPCWireRemote:input_type -> meshnet.v1beta1.WireDef + 8, // 14: meshnet.v1beta1.Remote.AddGRPCWiresRemoteBatch:input_type -> meshnet.v1beta1.WireDefBatch + 6, // 15: meshnet.v1beta1.Remote.GRPCWireDownRemote:input_type -> meshnet.v1beta1.WireDef + 11, // 16: meshnet.v1beta1.WireProtocol.SendToOnce:input_type -> meshnet.v1beta1.Packet + 11, // 17: meshnet.v1beta1.WireProtocol.SendToStream:input_type -> meshnet.v1beta1.Packet + 0, // 18: meshnet.v1beta1.Local.Get:output_type -> meshnet.v1beta1.Pod + 4, // 19: meshnet.v1beta1.Local.SetAlive:output_type -> meshnet.v1beta1.BoolResponse + 4, // 20: meshnet.v1beta1.Local.SkipReverse:output_type -> meshnet.v1beta1.BoolResponse + 4, // 21: meshnet.v1beta1.Local.Skip:output_type -> meshnet.v1beta1.BoolResponse + 4, // 22: meshnet.v1beta1.Local.IsSkipped:output_type -> meshnet.v1beta1.BoolResponse + 7, // 23: meshnet.v1beta1.Local.GRPCWireExists:output_type -> meshnet.v1beta1.WireCreateResponse + 4, // 24: meshnet.v1beta1.Local.AddGRPCWireLocal:output_type -> meshnet.v1beta1.BoolResponse + 4, // 25: meshnet.v1beta1.Local.RemGRPCWire:output_type -> meshnet.v1beta1.BoolResponse + 13, // 26: meshnet.v1beta1.Local.GenerateNodeInterfaceName:output_type -> meshnet.v1beta1.GenerateNodeInterfaceNameResponse + 4, // 27: meshnet.v1beta1.Remote.Update:output_type -> meshnet.v1beta1.BoolResponse + 7, // 28: meshnet.v1beta1.Remote.AddGRPCWireRemote:output_type -> meshnet.v1beta1.WireCreateResponse + 9, // 29: meshnet.v1beta1.Remote.AddGRPCWiresRemoteBatch:output_type -> meshnet.v1beta1.WireCreateResponseBatch + 10, // 30: meshnet.v1beta1.Remote.GRPCWireDownRemote:output_type -> meshnet.v1beta1.WireDownResponse + 4, // 31: meshnet.v1beta1.WireProtocol.SendToOnce:output_type -> meshnet.v1beta1.BoolResponse + 4, // 32: meshnet.v1beta1.WireProtocol.SendToStream:output_type -> meshnet.v1beta1.BoolResponse + 18, // [18:33] is the sub-list for method output_type + 3, // [3:18] is the sub-list for method input_type + 3, // [3:3] is the sub-list for extension type_name + 3, // [3:3] is the sub-list for extension extendee + 0, // [0:3] is the sub-list for field type_name } func init() { file_daemon_proto_meshnet_v1beta1_meshnet_proto_init() } @@ -1005,7 +1104,7 @@ func file_daemon_proto_meshnet_v1beta1_meshnet_proto_init() { GoPackagePath: reflect.TypeOf(x{}).PkgPath(), RawDescriptor: unsafe.Slice(unsafe.StringData(file_daemon_proto_meshnet_v1beta1_meshnet_proto_rawDesc), len(file_daemon_proto_meshnet_v1beta1_meshnet_proto_rawDesc)), NumEnums: 0, - NumMessages: 12, + NumMessages: 14, NumExtensions: 0, NumServices: 3, }, diff --git a/third_party/meshnet/daemon/proto/meshnet/v1beta1/meshnet.proto b/third_party/meshnet/daemon/proto/meshnet/v1beta1/meshnet.proto index 4d5b311d..ab5bc6dc 100644 --- a/third_party/meshnet/daemon/proto/meshnet/v1beta1/meshnet.proto +++ b/third_party/meshnet/daemon/proto/meshnet/v1beta1/meshnet.proto @@ -98,6 +98,14 @@ message WireCreateResponse { int64 peer_intf_id = 2; } +message WireDefBatch { + repeated WireDef items = 1; +} + +message WireCreateResponseBatch { + repeated WireCreateResponse items = 1; +} + message WireDownResponse { bool response = 1; } @@ -146,6 +154,7 @@ service Local { service Remote { rpc Update (RemotePod) returns (BoolResponse); rpc AddGRPCWireRemote(WireDef) returns (WireCreateResponse); + rpc AddGRPCWiresRemoteBatch(WireDefBatch) returns (WireCreateResponseBatch); rpc GRPCWireDownRemote(WireDef) returns (WireDownResponse); } diff --git a/third_party/meshnet/daemon/proto/meshnet/v1beta1/meshnet_grpc.pb.go b/third_party/meshnet/daemon/proto/meshnet/v1beta1/meshnet_grpc.pb.go index 17b8dad9..a2d60ee6 100644 --- a/third_party/meshnet/daemon/proto/meshnet/v1beta1/meshnet_grpc.pb.go +++ b/third_party/meshnet/daemon/proto/meshnet/v1beta1/meshnet_grpc.pb.go @@ -1,6 +1,6 @@ // Code generated by protoc-gen-go-grpc. DO NOT EDIT. // versions: -// - protoc-gen-go-grpc v1.6.0 +// - protoc-gen-go-grpc v1.6.1 // - protoc (unknown) // source: daemon/proto/meshnet/v1beta1/meshnet.proto @@ -440,9 +440,10 @@ var Local_ServiceDesc = grpc.ServiceDesc{ } const ( - Remote_Update_FullMethodName = "/meshnet.v1beta1.Remote/Update" - Remote_AddGRPCWireRemote_FullMethodName = "/meshnet.v1beta1.Remote/AddGRPCWireRemote" - Remote_GRPCWireDownRemote_FullMethodName = "/meshnet.v1beta1.Remote/GRPCWireDownRemote" + Remote_Update_FullMethodName = "/meshnet.v1beta1.Remote/Update" + Remote_AddGRPCWireRemote_FullMethodName = "/meshnet.v1beta1.Remote/AddGRPCWireRemote" + Remote_AddGRPCWiresRemoteBatch_FullMethodName = "/meshnet.v1beta1.Remote/AddGRPCWiresRemoteBatch" + Remote_GRPCWireDownRemote_FullMethodName = "/meshnet.v1beta1.Remote/GRPCWireDownRemote" ) // RemoteClient is the client API for Remote service. @@ -451,6 +452,7 @@ const ( type RemoteClient interface { Update(ctx context.Context, in *RemotePod, opts ...grpc.CallOption) (*BoolResponse, error) AddGRPCWireRemote(ctx context.Context, in *WireDef, opts ...grpc.CallOption) (*WireCreateResponse, error) + AddGRPCWiresRemoteBatch(ctx context.Context, in *WireDefBatch, opts ...grpc.CallOption) (*WireCreateResponseBatch, error) GRPCWireDownRemote(ctx context.Context, in *WireDef, opts ...grpc.CallOption) (*WireDownResponse, error) } @@ -482,6 +484,16 @@ func (c *remoteClient) AddGRPCWireRemote(ctx context.Context, in *WireDef, opts return out, nil } +func (c *remoteClient) AddGRPCWiresRemoteBatch(ctx context.Context, in *WireDefBatch, opts ...grpc.CallOption) (*WireCreateResponseBatch, error) { + cOpts := append([]grpc.CallOption{grpc.StaticMethod()}, opts...) + out := new(WireCreateResponseBatch) + err := c.cc.Invoke(ctx, Remote_AddGRPCWiresRemoteBatch_FullMethodName, in, out, cOpts...) + if err != nil { + return nil, err + } + return out, nil +} + func (c *remoteClient) GRPCWireDownRemote(ctx context.Context, in *WireDef, opts ...grpc.CallOption) (*WireDownResponse, error) { cOpts := append([]grpc.CallOption{grpc.StaticMethod()}, opts...) out := new(WireDownResponse) @@ -498,6 +510,7 @@ func (c *remoteClient) GRPCWireDownRemote(ctx context.Context, in *WireDef, opts type RemoteServer interface { Update(context.Context, *RemotePod) (*BoolResponse, error) AddGRPCWireRemote(context.Context, *WireDef) (*WireCreateResponse, error) + AddGRPCWiresRemoteBatch(context.Context, *WireDefBatch) (*WireCreateResponseBatch, error) GRPCWireDownRemote(context.Context, *WireDef) (*WireDownResponse, error) mustEmbedUnimplementedRemoteServer() } @@ -515,6 +528,9 @@ func (UnimplementedRemoteServer) Update(context.Context, *RemotePod) (*BoolRespo func (UnimplementedRemoteServer) AddGRPCWireRemote(context.Context, *WireDef) (*WireCreateResponse, error) { return nil, status.Error(codes.Unimplemented, "method AddGRPCWireRemote not implemented") } +func (UnimplementedRemoteServer) AddGRPCWiresRemoteBatch(context.Context, *WireDefBatch) (*WireCreateResponseBatch, error) { + return nil, status.Error(codes.Unimplemented, "method AddGRPCWiresRemoteBatch not implemented") +} func (UnimplementedRemoteServer) GRPCWireDownRemote(context.Context, *WireDef) (*WireDownResponse, error) { return nil, status.Error(codes.Unimplemented, "method GRPCWireDownRemote not implemented") } @@ -575,6 +591,24 @@ func _Remote_AddGRPCWireRemote_Handler(srv interface{}, ctx context.Context, dec return interceptor(ctx, in, info, handler) } +func _Remote_AddGRPCWiresRemoteBatch_Handler(srv interface{}, ctx context.Context, dec func(interface{}) error, interceptor grpc.UnaryServerInterceptor) (interface{}, error) { + in := new(WireDefBatch) + if err := dec(in); err != nil { + return nil, err + } + if interceptor == nil { + return srv.(RemoteServer).AddGRPCWiresRemoteBatch(ctx, in) + } + info := &grpc.UnaryServerInfo{ + Server: srv, + FullMethod: Remote_AddGRPCWiresRemoteBatch_FullMethodName, + } + handler := func(ctx context.Context, req interface{}) (interface{}, error) { + return srv.(RemoteServer).AddGRPCWiresRemoteBatch(ctx, req.(*WireDefBatch)) + } + return interceptor(ctx, in, info, handler) +} + func _Remote_GRPCWireDownRemote_Handler(srv interface{}, ctx context.Context, dec func(interface{}) error, interceptor grpc.UnaryServerInterceptor) (interface{}, error) { in := new(WireDef) if err := dec(in); err != nil { @@ -608,6 +642,10 @@ var Remote_ServiceDesc = grpc.ServiceDesc{ MethodName: "AddGRPCWireRemote", Handler: _Remote_AddGRPCWireRemote_Handler, }, + { + MethodName: "AddGRPCWiresRemoteBatch", + Handler: _Remote_AddGRPCWiresRemoteBatch_Handler, + }, { MethodName: "GRPCWireDownRemote", Handler: _Remote_GRPCWireDownRemote_Handler, diff --git a/third_party/meshnet/daemon/vxlan/vxlan.go b/third_party/meshnet/daemon/vxlan/vxlan.go index c1eaa300..9bf49c51 100644 --- a/third_party/meshnet/daemon/vxlan/vxlan.go +++ b/third_party/meshnet/daemon/vxlan/vxlan.go @@ -1,3 +1,4 @@ +// Package vxlan implements VXLAN overlay link creation and network interface management for meshnet daemon. package vxlan import ( @@ -12,10 +13,12 @@ import ( "github.com/vishvananda/netlink" mpb "github.com/openconfig/kne/third_party/meshnet/daemon/proto/meshnet/v1beta1" + "github.com/openconfig/kne/third_party/meshnet/utils/wireutil" ) var vxLanOvrlyLogger *log.Entry = nil +// InitLogger initializes the logrus logger for the VXLAN overlay daemon. func InitLogger() { vxLanOvrlyLogger = log.WithFields(log.Fields{"daemon": "meshnetd", "overlay": "vxLAN"}) } @@ -105,6 +108,20 @@ func CreateOrUpdate(v *mpb.RemotePod) error { } } + // Tune txqueuelen inside the container netns (configurable via LINK_TXQUEUELEN) + if podNs, err := ns.GetNS(veth.NsName); err == nil { + _ = podNs.Do(func(_ ns.NetNS) error { + if link, err := netlink.LinkByName(veth.LinkName); err == nil { + txqLen := wireutil.GetLinkTxQLen() + if err := netlink.LinkSetTxQLen(link, txqLen); err != nil { + vxLanOvrlyLogger.Warnf("failed to set txqueuelen %d on %s inside %s: %v", txqLen, veth.LinkName, veth.NsName, err) + } + } + return nil + }) + podNs.Close() + } + return nil } diff --git a/third_party/meshnet/docker/Dockerfile b/third_party/meshnet/docker/Dockerfile index 16b66d9b..ec2efd25 100644 --- a/third_party/meshnet/docker/Dockerfile +++ b/third_party/meshnet/docker/Dockerfile @@ -26,7 +26,6 @@ COPY go.sum . RUN go mod download && \ apt-get update -y && \ apt-get install -y --no-install-recommends \ - libpcap-dev \ libsystemd-dev && \ rm -rf /var/lib/apt/lists/* @@ -44,7 +43,7 @@ COPY --from=proto_base /src/ . RUN GOOS=${TARGETOS} GOARCH=${TARGETARCH} go build -o meshnet plugin/meshnet.go plugin/grpcwires-plugin.go && \ - GOOS=${TARGETOS} CGO_ENABLED=1 GOARCH=${TARGETARCH} go build -ldflags="-s -w" -o meshnetd daemon/main.go + GOOS=${TARGETOS} CGO_ENABLED=0 GOARCH=${TARGETARCH} go build -ldflags="-s -w" -o meshnetd daemon/main.go #----------------------------------------------------- Final Container --------------------------------- @@ -54,8 +53,7 @@ FROM debian:13-slim # hadolint ignore=DL3008 RUN apt-get update && \ apt-get install -y --no-install-recommends \ - jq \ - libpcap-dev && \ + jq && \ rm -rf /var/lib/apt/lists/* COPY --from=build /go/src/github.com/openconfig/kne/third_party/meshnet/meshnet / diff --git a/third_party/meshnet/plugin/grpcwires-plugin.go b/third_party/meshnet/plugin/grpcwires-plugin.go index ad8a4545..3de6f281 100644 --- a/third_party/meshnet/plugin/grpcwires-plugin.go +++ b/third_party/meshnet/plugin/grpcwires-plugin.go @@ -22,7 +22,7 @@ const ( skipStatusRetryCount = skipStatusRetryWarnCount * 4 // how many times to retry ) -// -------------------------------------------------------------------------------------------------------- +// CreatGRPCChan sets up the local and remote ends of a gRPC wire channel between two pods on different nodes. func CreatGRPCChan(link *mpb.Link, localPod *mpb.Pod, peerPod *mpb.Pod, localClient mpb.LocalClient, cniArgs *k8sArgs, ctx context.Context) error { // At this point pods attached to both end of this link are both up. They have got the management IP already. @@ -253,7 +253,7 @@ func CreatGRPCChan(link *mpb.Link, localPod *mpb.Pod, peerPod *mpb.Pod, localCli return nil } -// This function is called when a K8S pod is getting deleted. +// MakeGRPCChanDown signals the remote peer node to tear down the remote gRPC wire end when a pod is deleted. func MakeGRPCChanDown(link *mpb.Link, localPod *mpb.Pod, peerPod *mpb.Pod, ctx context.Context) error { if link == nil { return fmt.Errorf("can't remove remote grpc info. link not provided. link:%p", link) diff --git a/third_party/meshnet/plugin/meshnet.go b/third_party/meshnet/plugin/meshnet.go index 2eff874e..2aa99a07 100644 --- a/third_party/meshnet/plugin/meshnet.go +++ b/third_party/meshnet/plugin/meshnet.go @@ -208,32 +208,52 @@ func cmdAdd(args *skel.CmdArgs) error { log.Infof("Add[%s]: Successfully registered pod alive status with meshnet daemon", string(cniArgs.K8S_POD_NAME)) if len(localPod.Links) > 0 { - waitCtx, cancel := context.WithDeadline(ctx, startTime.Add(15*time.Second)) + waitCtx, cancel := context.WithDeadline(ctx, startTime.Add(30*time.Second)) defer cancel() - _ = ns.WithNetNSPath(args.Netns, func(_ ns.NetNS) error { - ticker := time.NewTicker(50 * time.Millisecond) - defer ticker.Stop() - for { - ready := true - for _, link := range localPod.Links { + + ticker := time.NewTicker(100 * time.Millisecond) + defer ticker.Stop() + + for { + allAreReady := true + for _, link := range localPod.Links { + // Check if interface exists in container netns + _ = ns.WithNetNSPath(args.Netns, func(_ ns.NetNS) error { if _, err := netlink.LinkByName(link.LocalIntf); err != nil { - ready = false - break + allAreReady = false } - } - if ready { - log.Infof("Add[%s]: All %d interfaces ready in container namespace", string(cniArgs.K8S_POD_NAME), len(localPod.Links)) return nil + }) + if !allAreReady { + break } - select { - case <-waitCtx.Done(): - log.Infof("Add[%s]: Bounded readiness wait expired (%d links); asynchronous completion will continue", string(cniArgs.K8S_POD_NAME), len(localPod.Links)) - return nil - case <-ticker.C: + + // For gRPC links, check if the gRPC wire is fully established on the daemon + wireDef := &mpb.WireDef{ + LocalPodNetNs: args.Netns, + LinkUid: link.Uid, + } + resp, err := meshnetClient.GRPCWireExists(waitCtx, wireDef) + if err != nil || !resp.Response { + allAreReady = false + break } } - }) + + if allAreReady { + log.Infof("Add[%s]: All %d interfaces and gRPC wires are ready", string(cniArgs.K8S_POD_NAME), len(localPod.Links)) + break + } + + select { + case <-waitCtx.Done(): + log.Warnf("Add[%s]: Readiness wait timed out (%d links); proceeding asynchronously", string(cniArgs.K8S_POD_NAME), len(localPod.Links)) + goto WaitDone + case <-ticker.C: + } + } } +WaitDone: return types.PrintResult(result, n.CNIVersion) } @@ -362,6 +382,7 @@ func cmdDel(args *skel.CmdArgs) error { return nil } +// SetInterNodeLinkType reads the inter-node link configuration file to set the default overlay mode (GRPC or VXLAN). func SetInterNodeLinkType() { // TODO: Find a more appropriate (if any) way to figure out intended link type // As of today, daemon gets the intended link type from env INTER_NODE_LINK_TYPE diff --git a/third_party/meshnet/utils/wireutil/sys_tune.go b/third_party/meshnet/utils/wireutil/sys_tune.go new file mode 100644 index 00000000..7e9c0ea7 --- /dev/null +++ b/third_party/meshnet/utils/wireutil/sys_tune.go @@ -0,0 +1,78 @@ +package wireutil + +import ( + "os" + "strconv" + + log "github.com/sirupsen/logrus" + "golang.org/x/sys/unix" +) + +// GetEnvInt reads an integer environment variable with a default fallback value if unset or invalid. +func GetEnvInt(key string, defaultVal int) int { + if valStr := os.Getenv(key); valStr != "" { + if val, err := strconv.Atoi(valStr); err == nil && val > 0 { + return val + } + } + return defaultVal +} + +func getEnvString(key string, defaultVal string) string { + if valStr := os.Getenv(key); valStr != "" { + return valStr + } + return defaultVal +} + +// GetLinkTxQLen returns the configured link txqueuelen (default 10000, configurable via LINK_TXQUEUELEN). +func GetLinkTxQLen() int { + return GetEnvInt("LINK_TXQUEUELEN", 10000) +} + +// TuneSystem configures global OS sysctl tunables (backlog, buffers, ARP/neighbor GC thresholds, +// multicast group limits, rp_filter, and IPv6 startup behavior) and RLIMIT_NOFILE for high-density, +// high-throughput network topologies. Values can be customized via environment variables. +func TuneSystem() { + // 1. Increase max open file descriptors for daemon (rlimit) + noFileLimit := GetEnvInt("RLIMIT_NOFILE", 1048576) + var rlim unix.Rlimit + rlim.Max = uint64(noFileLimit) + rlim.Cur = uint64(noFileLimit) + if err := unix.Setrlimit(unix.RLIMIT_NOFILE, &rlim); err != nil { + log.Warnf("TuneSystem: failed to set RLIMIT_NOFILE to %d: %v", noFileLimit, err) + } else { + log.Infof("TuneSystem: successfully set RLIMIT_NOFILE to %d", noFileLimit) + } + + // 2. Sysctl kernel tunables for network device backlog, buffer limits, ARP/neighbor GC thresholds, + // multicast memberships, reverse path filtering, and IPv6 DAD/RS startup tuning. + sysctls := map[string]string{ + "/proc/sys/net/core/netdev_max_backlog": getEnvString("NETDEV_MAX_BACKLOG", "10000"), + "/proc/sys/net/core/rmem_max": getEnvString("RMEM_MAX", "16777216"), + "/proc/sys/net/core/wmem_max": getEnvString("WMEM_MAX", "16777216"), + "/proc/sys/net/core/rmem_default": getEnvString("RMEM_DEFAULT", "16777216"), + "/proc/sys/net/core/wmem_default": getEnvString("WMEM_DEFAULT", "16777216"), + "/proc/sys/net/ipv4/neigh/default/gc_thresh1": getEnvString("ARP_GC_THRESH1", "1024"), + "/proc/sys/net/ipv4/neigh/default/gc_thresh2": getEnvString("ARP_GC_THRESH2", "4096"), + "/proc/sys/net/ipv4/neigh/default/gc_thresh3": getEnvString("ARP_GC_THRESH3", "8192"), + "/proc/sys/net/ipv6/neigh/default/gc_thresh1": getEnvString("ARP_GC_THRESH1", "1024"), + "/proc/sys/net/ipv6/neigh/default/gc_thresh2": getEnvString("ARP_GC_THRESH2", "4096"), + "/proc/sys/net/ipv6/neigh/default/gc_thresh3": getEnvString("ARP_GC_THRESH3", "8192"), + "/proc/sys/net/ipv4/igmp_max_memberships": getEnvString("IGMP_MAX_MEMBERSHIPS", "10000"), + "/proc/sys/net/ipv6/mld_max_msf": getEnvString("MLD_MAX_MSF", "4096"), + "/proc/sys/net/ipv4/conf/all/rp_filter": getEnvString("RP_FILTER", "2"), + "/proc/sys/net/ipv4/conf/default/rp_filter": getEnvString("RP_FILTER", "2"), + "/proc/sys/net/ipv6/conf/default/accept_dad": getEnvString("IPV6_ACCEPT_DAD", "0"), + "/proc/sys/net/ipv6/conf/default/router_solicitations": getEnvString("IPV6_ROUTER_SOLICITATIONS", "0"), + "/proc/sys/net/ipv6/route/max_size": getEnvString("IPV6_ROUTE_MAX_SIZE", "1048576"), + } + + for path, val := range sysctls { + if err := os.WriteFile(path, []byte(val), 0644); err != nil { + log.Warnf("TuneSystem: failed to write %s to %s: %v", val, path, err) + } else { + log.Infof("TuneSystem: set %s = %s", path, val) + } + } +} diff --git a/third_party/meshnet/utils/wireutil/tap.go b/third_party/meshnet/utils/wireutil/tap.go new file mode 100644 index 00000000..5cf12a41 --- /dev/null +++ b/third_party/meshnet/utils/wireutil/tap.go @@ -0,0 +1,88 @@ +package wireutil + +import ( + "fmt" + "os" + "unsafe" + + "github.com/containernetworking/plugins/pkg/ns" + log "github.com/sirupsen/logrus" + "github.com/vishvananda/netlink" + "golang.org/x/sys/unix" +) + +const tunDevice = "/dev/net/tun" + +type ifreq struct { + Name [unix.IFNAMSIZ]byte + Flags uint16 + _ [22]byte +} + +// CreateOrAttachTAP opens an existing persistent TAP device or creates a new persistent TAP device +// with the given ifName inside the specified network namespace at podNsPath. +// Returns the open *os.File handle to the TAP device. +func CreateOrAttachTAP(podNsPath string, ifName string, ipCIDR string) (*os.File, error) { + podNs, err := ns.GetNS(podNsPath) + if err != nil { + return nil, fmt.Errorf("could not open netns %s: %w", podNsPath, err) + } + defer podNs.Close() + + var tapFile *os.File + + err = podNs.Do(func(_ ns.NetNS) error { + fd, err := unix.Open(tunDevice, unix.O_RDWR, 0) + if err != nil { + return fmt.Errorf("failed to open %s in netns %s: %w", tunDevice, podNsPath, err) + } + + var ifr ifreq + copy(ifr.Name[:], []byte(ifName)) + ifr.Flags = unix.IFF_TAP | unix.IFF_NO_PI + + _, _, errno := unix.Syscall(unix.SYS_IOCTL, uintptr(fd), uintptr(unix.TUNSETIFF), uintptr(unsafe.Pointer(&ifr))) + if errno != 0 { + unix.Close(fd) + return fmt.Errorf("TUNSETIFF failed for %s in netns %s: %v", ifName, podNsPath, errno) + } + + // Make device persistent so it survives process crashes/restarts + _, _, _ = unix.Syscall(unix.SYS_IOCTL, uintptr(fd), uintptr(unix.TUNSETPERSIST), 1) + + link, err := netlink.LinkByName(ifName) + if err != nil { + unix.Close(fd) + return fmt.Errorf("failed to find link %s inside netns %s: %w", ifName, podNsPath, err) + } + + // Increase txqueuelen for high-throughput packet processing (configurable via LINK_TXQUEUELEN) + txqLen := GetLinkTxQLen() + if err := netlink.LinkSetTxQLen(link, txqLen); err != nil { + log.Warnf("CreateOrAttachTAP: failed to set txqueuelen %d on %s in netns %s: %v", txqLen, ifName, podNsPath, err) + } + + if err := netlink.LinkSetUp(link); err != nil { + log.Warnf("CreateOrAttachTAP: failed to set %s UP in netns %s: %v", ifName, podNsPath, err) + } + + if ipCIDR != "" { + addr, err := netlink.ParseAddr(ipCIDR) + if err == nil { + _ = netlink.AddrAdd(link, addr) + } + } + + tapFile = os.NewFile(uintptr(fd), ifName) + return nil + }) + + if err != nil { + return nil, err + } + + // Disable tx offload inside the netns (ignore error if non-fatal) + _ = SetTxChecksumOff(ifName, podNsPath) + + return tapFile, nil +} diff --git a/third_party/meshnet/utils/wireutil/veth.go b/third_party/meshnet/utils/wireutil/veth.go index 8add33e5..58ebfdb9 100644 --- a/third_party/meshnet/utils/wireutil/veth.go +++ b/third_party/meshnet/utils/wireutil/veth.go @@ -173,6 +173,12 @@ func ConfigurePodLinks(podNsPath string, links []PodLinkConfig) error { } } + // Increase txqueuelen for high-throughput packet processing (configurable via LINK_TXQUEUELEN) + txqLen := GetLinkTxQLen() + if err := netlink.LinkSetTxQLen(link, txqLen); err != nil { + log.Warnf("ConfigurePodLinks: failed to set txqueuelen %d on %s inside %s: %v", txqLen, cfg.LocalIntf, podNsPath, err) + } + if err := netlink.LinkSetUp(link); err != nil { return fmt.Errorf("failed to set %s UP inside %s: %w", cfg.LocalIntf, podNsPath, err) } diff --git a/third_party/meshnet/utils/wireutil/wire-util.go b/third_party/meshnet/utils/wireutil/wire-util.go index 9fda21a7..9994d51a 100644 --- a/third_party/meshnet/utils/wireutil/wire-util.go +++ b/third_party/meshnet/utils/wireutil/wire-util.go @@ -1,3 +1,5 @@ +// Package wireutil provides low-level network interface creation, TAP/veth management, +// checksum offload tuning, and OS performance utilities for meshnet. package wireutil import ( @@ -34,6 +36,8 @@ const ( INTER_NODE_LINK_GRPC = "GRPC" ) +// SetTxChecksumOff disables TX checksum and segmentation offloading on the specified interface +// inside the target network namespace to prevent checksum corruption during packet forwarding. func SetTxChecksumOff(intfName, nsName string) error { var vethNs ns.NetNS var err error