diff options
| -rw-r--r-- | common/constants/default.go | 3 | ||||
| -rw-r--r-- | edge/pkg/edged/edged.go | 11 | ||||
| -rw-r--r-- | edge/pkg/metamanager/client/pod.go | 17 | ||||
| -rw-r--r-- | edge/pkg/metamanager/dao/meta.go | 73 | ||||
| -rw-r--r-- | keadm/cmd/keadm/app/cmd/common/content.go | 8 | ||||
| -rw-r--r-- | keadm/cmd/keadm/app/cmd/edge/image.go | 3 | ||||
| -rw-r--r-- | keadm/cmd/keadm/app/cmd/edge/join.go | 10 | ||||
| -rw-r--r-- | keadm/cmd/keadm/app/cmd/edge/upgrade.go | 9 |
8 files changed, 122 insertions, 12 deletions
diff --git a/common/constants/default.go b/common/constants/default.go index fab9a5f16..72d852121 100644 --- a/common/constants/default.go +++ b/common/constants/default.go @@ -166,4 +166,7 @@ const ( EdgeNodeRoleKey = "node-role.kubernetes.io/edge" EdgeNodeRoleValue = "" + + DeafultMosquittoContainerName = "mqtt-kubeedge" + DeployMqttContainerEnv = "DEPLOY_MQTT_CONTAINER" ) diff --git a/edge/pkg/edged/edged.go b/edge/pkg/edged/edged.go index 3949ac0c2..172aefe05 100644 --- a/edge/pkg/edged/edged.go +++ b/edge/pkg/edged/edged.go @@ -29,6 +29,7 @@ import ( "encoding/json" "fmt" "os" + "strconv" "time" "github.com/container-storage-interface/spec/lib/go/csi" @@ -56,6 +57,7 @@ import ( kubebridge "github.com/kubeedge/kubeedge/edge/pkg/edged/kubeclientbridge" "github.com/kubeedge/kubeedge/edge/pkg/metamanager" metaclient "github.com/kubeedge/kubeedge/edge/pkg/metamanager/client" + "github.com/kubeedge/kubeedge/edge/pkg/metamanager/dao" "github.com/kubeedge/kubeedge/pkg/apis/componentconfig/edgecore/v1alpha2" "github.com/kubeedge/kubeedge/pkg/version" ) @@ -118,6 +120,15 @@ func (e *edged) Enable() bool { func (e *edged) Start() { klog.Info("Starting edged...") + // edged saves the data of mqtt container in sqlite3 and starts it. This is a temporary workaround and will be modified in v1.15. + withMqtt, err := strconv.ParseBool(os.Getenv(constants.DeployMqttContainerEnv)) + if err == nil && withMqtt { + err := dao.SaveMQTTMeta(e.nodeName) + if err != nil { + klog.ErrorS(err, "Start mqtt container failed") + } + } + go func() { err := DefaultRunLiteKubelet(e.context, e.KubeletServer, e.KubeletDeps, e.FeatureGate) if err != nil { diff --git a/edge/pkg/metamanager/client/pod.go b/edge/pkg/metamanager/client/pod.go index 8b4a9291f..fd3aafac0 100644 --- a/edge/pkg/metamanager/client/pod.go +++ b/edge/pkg/metamanager/client/pod.go @@ -98,6 +98,9 @@ func (c *pods) Get(name string) (*corev1.Pod, error) { func (c *pods) Patch(name string, patchBytes []byte) (*corev1.Pod, error) { resource := fmt.Sprintf("%s/%s/%s", c.namespace, model.ResourceTypePodPatch, name) + if name == constants.DeafultMosquittoContainerName { + return handleMqttMeta() + } podMsg := message.BuildMsg(modules.MetaGroup, "", modules.EdgedModuleName, resource, model.PatchOperation, string(patchBytes)) resp, err := c.send.SendSync(podMsg) if err != nil { @@ -173,3 +176,17 @@ func updatePodDB(resource string, pod *corev1.Pod) error { Value: string(podContent)} return dao.InsertOrUpdate(meta) } + +func handleMqttMeta() (*corev1.Pod, error) { + var pod corev1.Pod + metas, err := dao.QueryMeta("key", fmt.Sprintf("default/pod/%s", constants.DeafultMosquittoContainerName)) + if err != nil || len(*metas) != 1 { + return nil, fmt.Errorf("get mqtt meta failed, err: %v", err) + } + + err = json.Unmarshal([]byte((*metas)[0]), &pod) + if err != nil { + return nil, fmt.Errorf("unmarshal mqtt meta failed, err: %v", err) + } + return &pod, nil +} diff --git a/edge/pkg/metamanager/dao/meta.go b/edge/pkg/metamanager/dao/meta.go index 09a6179a0..93d1ad8ce 100644 --- a/edge/pkg/metamanager/dao/meta.go +++ b/edge/pkg/metamanager/dao/meta.go @@ -1,11 +1,16 @@ package dao import ( + "encoding/json" "fmt" "strings" + coreV1 "k8s.io/api/core/v1" + metaV1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/util/uuid" "k8s.io/klog/v2" + "github.com/kubeedge/kubeedge/common/constants" "github.com/kubeedge/kubeedge/edge/pkg/common/dbm" ) @@ -112,3 +117,71 @@ func QueryAllMeta(key string, condition string) (*[]Meta, error) { return meta, nil } + +// SaveMQTTMeta saves mqtt container data in sqlites +// When egdecore starts, edged will start mqtt container +func SaveMQTTMeta(nodeName string) error { + flag := true + mqttData := coreV1.Pod{ + TypeMeta: metaV1.TypeMeta{ + Kind: "Pod", + APIVersion: "v1", + }, + ObjectMeta: metaV1.ObjectMeta{ + Name: constants.DeafultMosquittoContainerName, + Namespace: "default", + UID: uuid.NewUUID(), + }, + Spec: coreV1.PodSpec{ + Containers: []coreV1.Container{ + { + Name: "mqtt", + Image: constants.DefaultMosquittoImage, + Ports: []coreV1.ContainerPort{ + { + ContainerPort: 1883, + HostPort: 1883, + Protocol: coreV1.ProtocolTCP, + }, { + ContainerPort: 9001, + HostPort: 9001, + Protocol: coreV1.ProtocolTCP, + }, + }, + VolumeMounts: []coreV1.VolumeMount{ + { + MountPath: "/mosquitto", + Name: "mqtt-path", + }, + }, + }, + }, + Volumes: []coreV1.Volume{ + { + Name: "mqtt-path", + VolumeSource: coreV1.VolumeSource{ + HostPath: &coreV1.HostPathVolumeSource{ + Path: "/var/lib/kubeedge/mqtt", + }, + }, + }, + }, + NodeName: nodeName, + RestartPolicy: coreV1.RestartPolicyAlways, + DNSPolicy: coreV1.DNSClusterFirst, + EnableServiceLinks: &flag, + }, + } + mqttDataStr, _ := json.Marshal(mqttData) + mqttMeta := Meta{ + Key: fmt.Sprintf("default/pod/%s", constants.DeafultMosquittoContainerName), + Type: "pod", + Value: string(mqttDataStr), + } + err := SaveMeta(&mqttMeta) + if err != nil { + return err + } + + return nil +} diff --git a/keadm/cmd/keadm/app/cmd/common/content.go b/keadm/cmd/keadm/app/cmd/common/content.go index ff9153d93..18db7682c 100644 --- a/keadm/cmd/keadm/app/cmd/common/content.go +++ b/keadm/cmd/keadm/app/cmd/common/content.go @@ -19,6 +19,8 @@ package common import ( "fmt" "os" + + "github.com/kubeedge/kubeedge/common/constants" ) // TODO (@zc2638) Need to migrate util's constants to common @@ -31,14 +33,16 @@ Type=simple ExecStart=%s Restart=always RestartSec=10 +Environment=%s [Install] WantedBy=multi-user.target ` -func GenerateServiceFile(process string, execStartCmd string) error { +func GenerateServiceFile(process string, execStartCmd string, withMqtt bool) error { filename := fmt.Sprintf("%s.service", process) - content := fmt.Sprintf(serviceFileTemplate, process, execStartCmd) + + content := fmt.Sprintf(serviceFileTemplate, process, execStartCmd, fmt.Sprintf("%s=%t", constants.DeployMqttContainerEnv, withMqtt)) serviceFilePath := fmt.Sprintf("/etc/systemd/system/%s", filename) return os.WriteFile(serviceFilePath, []byte(content), os.ModePerm) } diff --git a/keadm/cmd/keadm/app/cmd/edge/image.go b/keadm/cmd/keadm/app/cmd/edge/image.go index 71e16efc1..a62902dee 100644 --- a/keadm/cmd/keadm/app/cmd/edge/image.go +++ b/keadm/cmd/keadm/app/cmd/edge/image.go @@ -53,9 +53,6 @@ func request(opt *common.JoinOptions, step *common.Step) error { if err := createMQTTConfigFile(); err != nil { return fmt.Errorf("create MQTT config file failed: %v", err) } - if err := runtime.RunMQTT(imageSet.Get(image.EdgeMQTT)); err != nil { - return fmt.Errorf("run MQTT failed: %v", err) - } } return nil } diff --git a/keadm/cmd/keadm/app/cmd/edge/join.go b/keadm/cmd/keadm/app/cmd/edge/join.go index 0053b8efb..5435ee5c7 100644 --- a/keadm/cmd/keadm/app/cmd/edge/join.go +++ b/keadm/cmd/keadm/app/cmd/edge/join.go @@ -186,7 +186,7 @@ func join(opt *common.JoinOptions, step *common.Step) error { } step.Printf("Generate systemd service file") - if err := common.GenerateServiceFile(util.KubeEdgeBinaryName, filepath.Join(util.KubeEdgeUsrBinPath, util.KubeEdgeBinaryName)); err != nil { + if err := common.GenerateServiceFile(util.KubeEdgeBinaryName, filepath.Join(util.KubeEdgeUsrBinPath, util.KubeEdgeBinaryName), opt.WithMQTT); err != nil { return fmt.Errorf("create systemd service file failed: %v", err) } @@ -203,7 +203,7 @@ func join(opt *common.JoinOptions, step *common.Step) error { } step.Printf("Run EdgeCore daemon") - err := runEdgeCore() + err := runEdgeCore(opt.WithMQTT) if err != nil { return fmt.Errorf("start edgecore failed: %v", err) } @@ -405,7 +405,7 @@ func setEdgedNodeLabels(opt *common.JoinOptions) map[string]string { return labelsMap } -func runEdgeCore() error { +func runEdgeCore(withMqtt bool) error { systemdExist := util.HasSystemd() var binExec, tip string @@ -416,6 +416,10 @@ func runEdgeCore() error { common.EdgeCore, common.EdgeCore) } else { tip = fmt.Sprintf("KubeEdge edgecore is running, For logs visit: %s%s.log", util.KubeEdgeLogPath, util.KubeEdgeBinaryName) + err := os.Setenv(constants.DeployMqttContainerEnv, strconv.FormatBool(withMqtt)) + if err != nil { + klog.Errorf("Set Environment %s failed, err: %v", constants.DeployMqttContainerEnv, err) + } binExec = fmt.Sprintf( "%s > %skubeedge/edge/%s.log 2>&1 &", filepath.Join(util.KubeEdgeUsrBinPath, util.KubeEdgeBinaryName), diff --git a/keadm/cmd/keadm/app/cmd/edge/upgrade.go b/keadm/cmd/keadm/app/cmd/edge/upgrade.go index 576d9d99c..4cc59f592 100644 --- a/keadm/cmd/keadm/app/cmd/edge/upgrade.go +++ b/keadm/cmd/keadm/app/cmd/edge/upgrade.go @@ -243,16 +243,17 @@ func (up *Upgrade) Process() error { return fmt.Errorf("failed to cp file: %v", err) } + // set withMqtt to false during upgrading edgecore, it will not affect the MQTT container. This is a temporary workaround and will be modified in v1.15. // generate edgecore.service if util.HasSystemd() { - err = common.GenerateServiceFile(util.KubeEdgeBinaryName, fmt.Sprintf("%s --config %s", filepath.Join(util.KubeEdgeUsrBinPath, util.KubeEdgeBinaryName), up.ConfigFilePath)) + err = common.GenerateServiceFile(util.KubeEdgeBinaryName, fmt.Sprintf("%s --config %s", filepath.Join(util.KubeEdgeUsrBinPath, util.KubeEdgeBinaryName), up.ConfigFilePath), false) if err != nil { return fmt.Errorf("failed to create edgecore.service file: %v", err) } } // start new edgecore service - err = runEdgeCore() + err = runEdgeCore(false) if err != nil { return fmt.Errorf("failed to start edgecore: %v", err) } @@ -287,14 +288,14 @@ func (up *Upgrade) Rollback() error { // generate edgecore.service if util.HasSystemd() { - err = common.GenerateServiceFile(util.KubeEdgeBinaryName, fmt.Sprintf("%s --config %s", filepath.Join(util.KubeEdgeUsrBinPath, util.KubeEdgeBinaryName), up.ConfigFilePath)) + err = common.GenerateServiceFile(util.KubeEdgeBinaryName, fmt.Sprintf("%s --config %s", filepath.Join(util.KubeEdgeUsrBinPath, util.KubeEdgeBinaryName), up.ConfigFilePath), false) if err != nil { return fmt.Errorf("failed to create edgecore.service file: %v", err) } } // start edgecore - err = runEdgeCore() + err = runEdgeCore(false) if err != nil { return fmt.Errorf("failed to start origin edgecore: %v", err) } |
