summaryrefslogtreecommitdiff
path: root/edge
diff options
context:
space:
mode:
authorShelley-BaoYue <baoyue2@huawei.com>2023-06-06 20:13:40 +0800
committerShelley-BaoYue <baoyue2@huawei.com>2023-07-24 10:27:10 +0800
commitc87f74756d11b1203c42a7fde382611b74a38531 (patch)
treef4149ddab370f6b46dd9c160158382c715b2c2f0 /edge
parentMerge pull request #4696 from Shelley-BaoYue/save-patchpod-meta (diff)
downloadkubeedge-c87f74756d11b1203c42a7fde382611b74a38531.tar.gz
fix mqtt: start mqtt container
Signed-off-by: Shelley-BaoYue <baoyue2@huawei.com>
Diffstat (limited to 'edge')
-rw-r--r--edge/cmd/edgecore/app/server.go2
-rw-r--r--edge/cmd/edgemark/hollow_edgecore.go2
-rw-r--r--edge/pkg/edged/edged.go15
-rw-r--r--edge/pkg/metamanager/client/pod.go17
-rw-r--r--edge/pkg/metamanager/dao/meta.go73
-rw-r--r--edge/test/integration/utils/test_setup.go1
6 files changed, 105 insertions, 5 deletions
diff --git a/edge/cmd/edgecore/app/server.go b/edge/cmd/edgecore/app/server.go
index 0438829f1..2d874ef1c 100644
--- a/edge/cmd/edgecore/app/server.go
+++ b/edge/cmd/edgecore/app/server.go
@@ -192,7 +192,7 @@ func environmentCheck() error {
// registerModules register all the modules started in edgecore
func registerModules(c *v1alpha2.EdgeCoreConfig) {
devicetwin.Register(c.Modules.DeviceTwin, c.Modules.Edged.HostnameOverride)
- edged.Register(c.Modules.Edged)
+ edged.Register(c.Modules.Edged, c.Modules.EventBus.EnableMqttContainer)
edgehub.Register(c.Modules.EdgeHub, c.Modules.Edged.HostnameOverride)
eventbus.Register(c.Modules.EventBus, c.Modules.Edged.HostnameOverride)
metamanager.Register(c.Modules.MetaManager)
diff --git a/edge/cmd/edgemark/hollow_edgecore.go b/edge/cmd/edgemark/hollow_edgecore.go
index 08de9d65f..a7cf1d6ce 100644
--- a/edge/cmd/edgemark/hollow_edgecore.go
+++ b/edge/cmd/edgemark/hollow_edgecore.go
@@ -112,7 +112,7 @@ func run(config *hollowEdgeNodeConfig) {
return kubeletapp.RunKubelet(s, kubeDeps, false)
}
- edged.Register(c.Modules.Edged)
+ edged.Register(c.Modules.Edged, c.Modules.EventBus.EnableMqttContainer)
edgehub.Register(c.Modules.EdgeHub, c.Modules.Edged.HostnameOverride)
metamanager.Register(c.Modules.MetaManager)
diff --git a/edge/pkg/edged/edged.go b/edge/pkg/edged/edged.go
index 8f73022fe..e1897b6f6 100644
--- a/edge/pkg/edged/edged.go
+++ b/edge/pkg/edged/edged.go
@@ -57,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"
)
@@ -88,14 +89,15 @@ type edged struct {
context context.Context
nodeName string
namespace string
+ withMqtt bool
}
var _ core.Module = (*edged)(nil)
// Register register edged
-func Register(e *v1alpha2.Edged) {
+func Register(e *v1alpha2.Edged, withMqtt bool) {
edgedconfig.InitConfigure(e)
- edged, err := newEdged(e.Enable, e.HostnameOverride, e.RegisterNodeNamespace)
+ edged, err := newEdged(e.Enable, e.HostnameOverride, e.RegisterNodeNamespace, withMqtt)
if err != nil {
klog.Errorf("init new edged error, %v", err)
os.Exit(1)
@@ -119,6 +121,11 @@ func (e *edged) Enable() bool {
func (e *edged) Start() {
klog.Info("Starting edged...")
+ 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 {
@@ -130,7 +137,7 @@ func (e *edged) Start() {
}
// newEdged creates new edged object and initialises it
-func newEdged(enable bool, nodeName, namespace string) (*edged, error) {
+func newEdged(enable bool, nodeName, namespace string, withMqtt bool) (*edged, error) {
var ed *edged
var err error
if !enable {
@@ -138,6 +145,7 @@ func newEdged(enable bool, nodeName, namespace string) (*edged, error) {
enable: enable,
nodeName: nodeName,
namespace: namespace,
+ withMqtt: withMqtt,
}, nil
}
@@ -179,6 +187,7 @@ func newEdged(enable bool, nodeName, namespace string) (*edged, error) {
FeatureGate: utilfeature.DefaultFeatureGate,
nodeName: nodeName,
namespace: namespace,
+ withMqtt: withMqtt,
}
return ed, nil
diff --git a/edge/pkg/metamanager/client/pod.go b/edge/pkg/metamanager/client/pod.go
index 2df2288bd..b859f22e5 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.DeafultMosquittoContianerName {
+ 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.DeafultMosquittoContianerName))
+ 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..53f1bcb38 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.DeafultMosquittoContianerName,
+ 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.DeafultMosquittoContianerName),
+ Type: "pod",
+ Value: string(mqttDataStr),
+ }
+ err := SaveMeta(&mqttMeta)
+ if err != nil {
+ return err
+ }
+
+ return nil
+}
diff --git a/edge/test/integration/utils/test_setup.go b/edge/test/integration/utils/test_setup.go
index 8890362f0..714485500 100644
--- a/edge/test/integration/utils/test_setup.go
+++ b/edge/test/integration/utils/test_setup.go
@@ -41,6 +41,7 @@ func CreateEdgeCoreConfigFile(nodeName string) error {
c.Modules.EdgeHub.TLSCertFile = "/tmp/edgecore/kubeedge.crt"
c.Modules.EdgeHub.TLSPrivateKeyFile = "/tmp/edgecore/kubeedge.key"
c.Modules.EventBus.Enable = true
+ c.Modules.EventBus.EnableMqttContainer = true
c.Modules.EventBus.MqttMode = edgecore.MqttModeInternal
c.Modules.DBTest.Enable = true
c.DataBase.DataSource = DBFile