summaryrefslogtreecommitdiff
diff options
context:
space:
mode:
authorKubeEdge Bot <48982446+kubeedge-bot@users.noreply.github.com>2023-07-27 21:49:51 +0800
committerGitHub <noreply@github.com>2023-07-27 21:49:51 +0800
commit02bb0926f5f63f9f5d9e84b1da4f879543174992 (patch)
tree94c5a5695c473dd99344145aa10febe7e303a4c8
parentMerge pull request #4901 from Shelley-BaoYue/automated-cherry-pick-of-#4888-u... (diff)
parentfix start mqtt container failed (diff)
downloadkubeedge-02bb0926f5f63f9f5d9e84b1da4f879543174992.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.go3
-rw-r--r--edge/pkg/edged/edged.go14
-rw-r--r--edge/pkg/metamanager/client/podstatus.go4
-rw-r--r--edge/pkg/metamanager/dao/meta.go74
-rw-r--r--keadm/cmd/keadm/app/cmd/common/content.go7
-rw-r--r--keadm/cmd/keadm/app/cmd/edge/image.go3
-rw-r--r--keadm/cmd/keadm/app/cmd/edge/join.go10
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),