diff options
| author | KubeEdge Bot <48982446+kubeedge-bot@users.noreply.github.com> | 2023-07-27 21:49:51 +0800 |
|---|---|---|
| committer | GitHub <noreply@github.com> | 2023-07-27 21:49:51 +0800 |
| commit | 02bb0926f5f63f9f5d9e84b1da4f879543174992 (patch) | |
| tree | 94c5a5695c473dd99344145aa10febe7e303a4c8 | |
| parent | Merge pull request #4901 from Shelley-BaoYue/automated-cherry-pick-of-#4888-u... (diff) | |
| parent | fix start mqtt container failed (diff) | |
| download | kubeedge-origin/release-1.11.tar.gz | |
Merge pull request #4900 from Shelley-BaoYue/cherrypick-4870-1.11v1.11.3origin/release-1.11
Cherry pick of #4814: fix start mqtt container failed
| -rw-r--r-- | common/constants/default.go | 3 | ||||
| -rw-r--r-- | edge/pkg/edged/edged.go | 14 | ||||
| -rw-r--r-- | edge/pkg/metamanager/client/podstatus.go | 4 | ||||
| -rw-r--r-- | edge/pkg/metamanager/dao/meta.go | 74 | ||||
| -rw-r--r-- | keadm/cmd/keadm/app/cmd/common/content.go | 7 | ||||
| -rw-r--r-- | keadm/cmd/keadm/app/cmd/edge/image.go | 3 | ||||
| -rw-r--r-- | keadm/cmd/keadm/app/cmd/edge/join.go | 10 |
7 files changed, 106 insertions, 9 deletions
diff --git a/common/constants/default.go b/common/constants/default.go index 75d686849..408435309 100644 --- a/common/constants/default.go +++ b/common/constants/default.go @@ -148,4 +148,7 @@ const ( DefaultBurst = 60 // MaxRespBodyLength is the max length of http response body MaxRespBodyLength = 1 << 20 // 1 MiB + + DeafultMosquittoContainerName = "mqtt-kubeedge" + DeployMqttContainerEnv = "DEPLOY_MQTT_CONTAINER" ) diff --git a/edge/pkg/edged/edged.go b/edge/pkg/edged/edged.go index a2fb7b763..2bbe6ab06 100644 --- a/edge/pkg/edged/edged.go +++ b/edge/pkg/edged/edged.go @@ -32,6 +32,7 @@ import ( "os" "path" "runtime" + "strconv" "strings" "sync" "time" @@ -117,6 +118,7 @@ import ( csiplugin "github.com/kubeedge/kubeedge/edge/pkg/edged/volume/csi" "github.com/kubeedge/kubeedge/edge/pkg/metamanager" "github.com/kubeedge/kubeedge/edge/pkg/metamanager/client" + "github.com/kubeedge/kubeedge/edge/pkg/metamanager/dao" "github.com/kubeedge/kubeedge/pkg/apis/componentconfig/edgecore/v1alpha1" "github.com/kubeedge/kubeedge/pkg/version" ) @@ -302,6 +304,16 @@ func (e *edged) GetRequestedContainersInfo(containerName string, options cadviso 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") + } + } + e.volumePluginMgr = NewInitializedVolumePluginMgr(e, ProbeVolumePlugins("")) if err := e.initializeModules(); err != nil { @@ -371,7 +383,7 @@ func (e *edged) Start() { go e.pluginManager.Run(edgedutil.NewSourcesReady(e.isInitPodReady), utilwait.NeverStop) // start the CPU manager in the clcm - err := e.clcm.StartCPUManager(e.GetActivePods, edgedutil.NewSourcesReady(e.isInitPodReady), e.statusManager, e.runtimeService) + err = e.clcm.StartCPUManager(e.GetActivePods, edgedutil.NewSourcesReady(e.isInitPodReady), e.statusManager, e.runtimeService) if err != nil { klog.Errorf("Failed to start container manager, err: %v", err) return diff --git a/edge/pkg/metamanager/client/podstatus.go b/edge/pkg/metamanager/client/podstatus.go index c6ad198af..3e4c275de 100644 --- a/edge/pkg/metamanager/client/podstatus.go +++ b/edge/pkg/metamanager/client/podstatus.go @@ -40,6 +40,10 @@ func (c *podStatus) Create(ps *edgeapi.PodStatusRequest) (*edgeapi.PodStatusRequ } func (c *podStatus) Update(rsName string, ps edgeapi.PodStatusRequest) error { + if ps.Name == constants.DeafultMosquittoContainerName { + return nil + } + podStatusMsg := message.BuildMsg(commodule.MetaGroup, "", commodule.EdgedModuleName, c.namespace+"/"+model.ResourceTypePodStatus+"/"+rsName, model.UpdateOperation, ps) resp, err := c.send.SendSync(podStatusMsg) if err != nil { diff --git a/edge/pkg/metamanager/dao/meta.go b/edge/pkg/metamanager/dao/meta.go index 07eb720c7..5796cf4d5 100644 --- a/edge/pkg/metamanager/dao/meta.go +++ b/edge/pkg/metamanager/dao/meta.go @@ -1,10 +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" ) @@ -100,3 +106,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: "eclipse-mosquitto:1.6.15", + 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 1660449ae..a42104c22 100644 --- a/keadm/cmd/keadm/app/cmd/common/content.go +++ b/keadm/cmd/keadm/app/cmd/common/content.go @@ -20,6 +20,8 @@ import ( "fmt" "os" "path/filepath" + + "github.com/kubeedge/kubeedge/common/constants" ) // TODO (@zc2638) Need to migrate util's constants to common @@ -32,14 +34,15 @@ Type=simple ExecStart=%s Restart=always RestartSec=10 +Environment=%s [Install] WantedBy=multi-user.target ` -func GenerateServiceFile(processPath string, process string) error { +func GenerateServiceFile(processPath string, process string, withMqtt bool) error { filename := fmt.Sprintf("%s.service", process) - content := fmt.Sprintf(serviceFileTemplate, process, filepath.Join(processPath, process)) + content := fmt.Sprintf(serviceFileTemplate, process, filepath.Join(processPath, process), 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 ba09fc96b..c9d229e7c 100644 --- a/keadm/cmd/keadm/app/cmd/edge/image.go +++ b/keadm/cmd/keadm/app/cmd/edge/image.go @@ -56,9 +56,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 4bc918d32..0b76aeb22 100644 --- a/keadm/cmd/keadm/app/cmd/edge/join.go +++ b/keadm/cmd/keadm/app/cmd/edge/join.go @@ -178,7 +178,7 @@ func join(opt *common.JoinOptions, step *common.Step) error { } step.Printf("Generate systemd service file") - if err := common.GenerateServiceFile(util.KubeEdgeUsrBinPath, util.KubeEdgeBinaryName); err != nil { + if err := common.GenerateServiceFile(util.KubeEdgeUsrBinPath, util.KubeEdgeBinaryName, opt.WithMQTT); err != nil { return fmt.Errorf("create systemd service file failed: %v", err) } @@ -188,7 +188,7 @@ func join(opt *common.JoinOptions, step *common.Step) error { } step.Printf("Run EdgeCore daemon") - return runEdgeCore() + return runEdgeCore(opt.WithMQTT) } func createDirs() error { @@ -291,7 +291,7 @@ func createEdgeConfigFiles(opt *common.JoinOptions) error { return common.Write2File(configFilePath, edgeCoreConfig) } -func runEdgeCore() error { +func runEdgeCore(withMqtt bool) error { systemdExist := util.HasSystemd() var binExec, tip string @@ -302,6 +302,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), |
