summaryrefslogtreecommitdiff
diff options
context:
space:
mode:
-rw-r--r--common/constants/default.go3
-rw-r--r--edge/pkg/edged/edged.go11
-rw-r--r--edge/pkg/metamanager/client/pod.go17
-rw-r--r--edge/pkg/metamanager/dao/meta.go73
-rw-r--r--keadm/cmd/keadm/app/cmd/common/content.go8
-rw-r--r--keadm/cmd/keadm/app/cmd/edge/image.go3
-rw-r--r--keadm/cmd/keadm/app/cmd/edge/join.go10
-rw-r--r--keadm/cmd/keadm/app/cmd/edge/upgrade.go9
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)
}