diff options
| author | Shelley-BaoYue <baoyue2@huawei.com> | 2023-06-06 20:13:40 +0800 |
|---|---|---|
| committer | Shelley-BaoYue <baoyue2@huawei.com> | 2023-07-24 10:27:10 +0800 |
| commit | c87f74756d11b1203c42a7fde382611b74a38531 (patch) | |
| tree | f4149ddab370f6b46dd9c160158382c715b2c2f0 /edge | |
| parent | Merge pull request #4696 from Shelley-BaoYue/save-patchpod-meta (diff) | |
| download | kubeedge-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.go | 2 | ||||
| -rw-r--r-- | edge/cmd/edgemark/hollow_edgecore.go | 2 | ||||
| -rw-r--r-- | edge/pkg/edged/edged.go | 15 | ||||
| -rw-r--r-- | edge/pkg/metamanager/client/pod.go | 17 | ||||
| -rw-r--r-- | edge/pkg/metamanager/dao/meta.go | 73 | ||||
| -rw-r--r-- | edge/test/integration/utils/test_setup.go | 1 |
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 |
