diff options
| author | luomengY <2938893385@qq.com> | 2023-08-30 21:50:01 +0800 |
|---|---|---|
| committer | luomengY <2938893385@qq.com> | 2023-09-23 15:41:03 +0800 |
| commit | 146654ade70f96302a45a75988ff8e32cfa51217 (patch) | |
| tree | 6208f3c309883586abe5ae5d1ea48f18fe5bf7c4 /cloud | |
| parent | Merge pull request #5006 from wlq1212/k8scompatibility/schedule (diff) | |
| download | kubeedge-146654ade70f96302a45a75988ff8e32cfa51217.tar.gz | |
add support static in kubeedge
Signed-off-by: luomengY <2938893385@qq.com>
Diffstat (limited to 'cloud')
| -rw-r--r-- | cloud/pkg/cloudhub/dispatcher/message_dispatcher.go | 3 | ||||
| -rw-r--r-- | cloud/pkg/edgecontroller/controller/upstream.go | 56 |
2 files changed, 58 insertions, 1 deletions
diff --git a/cloud/pkg/cloudhub/dispatcher/message_dispatcher.go b/cloud/pkg/cloudhub/dispatcher/message_dispatcher.go index cd71bacdd..8dc837c7c 100644 --- a/cloud/pkg/cloudhub/dispatcher/message_dispatcher.go +++ b/cloud/pkg/cloudhub/dispatcher/message_dispatcher.go @@ -467,7 +467,8 @@ func noAckRequired(msg *beehivemodel.Message) bool { resourceType == beehivemodel.ResourceTypeLease || resourceType == beehivemodel.ResourceTypeNodePatch || resourceType == beehivemodel.ResourceTypePodPatch || - resourceType == beehivemodel.ResourceTypePodStatus { + resourceType == beehivemodel.ResourceTypePodStatus || + resourceType == beehivemodel.ResourceTypeCreatePod { return true } } diff --git a/cloud/pkg/edgecontroller/controller/upstream.go b/cloud/pkg/edgecontroller/controller/upstream.go index 8e38363a8..7226d2922 100644 --- a/cloud/pkg/edgecontroller/controller/upstream.go +++ b/cloud/pkg/edgecontroller/controller/upstream.go @@ -116,6 +116,7 @@ type UpstreamController struct { ruleStatusChan chan model.Message createLeaseChan chan model.Message queryLeaseChan chan model.Message + createPodChan chan model.Message // lister podLister corelisters.PodLister @@ -182,6 +183,9 @@ func (uc *UpstreamController) Start() error { for i := 0; i < int(uc.config.Load.UpdateRuleStatusWorkers); i++ { go uc.updateRuleStatus() } + for i := 0; i < int(uc.config.Load.CreatePodWorks); i++ { + go uc.createPod() + } return nil } @@ -257,6 +261,9 @@ func (uc *UpstreamController) dispatchMessage() { case model.QueryOperation: uc.queryLeaseChan <- msg } + case model.ResourceTypeCreatePod: + uc.createPodChan <- msg + default: klog.Errorf("message: %s, resource type: %s unsupported", msg.GetID(), resourceType) } @@ -1042,6 +1049,54 @@ func (uc *UpstreamController) patchPod() { } } +func (uc *UpstreamController) createPod() { + for { + select { + case <-beehiveContext.Done(): + klog.Warning("stop createPod") + return + case msg := <-uc.createPodChan: + klog.V(5).Infof("message: %s, operation is: %s, and resource is %s", msg.GetID(), msg.GetOperation(), msg.GetResource()) + namespace, err := messagelayer.GetNamespace(msg) + if err != nil { + klog.Warningf("message: %s process failure, get namespace failed with error: %v", msg.GetID(), err) + continue + } + name, err := messagelayer.GetResourceName(msg) + if err != nil { + klog.Warningf("message: %s process failure, get resource name failed with error: %v", msg.GetID(), err) + continue + } + + podBytes, err := msg.GetContentData() + if err != nil { + klog.Warningf("message: %s process failure, get data failed with error: %v", msg.GetID(), err) + continue + } + var pod v1.Pod + if err = json.Unmarshal(podBytes, &pod); err != nil { + klog.Errorf("unmarshal pod request failed with error: %v", err) + continue + } + + createPod, err := uc.kubeClient.CoreV1().Pods(namespace).Create(context.TODO(), &pod, metaV1.CreateOptions{}) + if err != nil { + klog.Errorf("message: %s process failure, create pod failed with error: %v, namespace: %s, name: %s", msg.GetID(), err, namespace, name) + } + + resMsg := model.NewMessage(msg.GetID()). + SetResourceVersion(createPod.ResourceVersion). + FillBody(&edgeapi.ObjectResp{Object: createPod, Err: err}). + BuildRouter(modules.EdgeControllerModuleName, constants.GroupResource, msg.GetResource(), model.ResponseOperation) + if err = uc.messageLayer.Response(*resMsg); err != nil { + klog.Errorf("Message: %s process failure, response failed with error: %v", msg.GetID(), err) + continue + } + klog.V(4).Infof("message: %s, create pod successfully, namespace: %s, name: %s", msg.GetID(), namespace, name) + } + } +} + func (uc *UpstreamController) deletePod() { for { select { @@ -1378,6 +1433,7 @@ func NewUpstreamController(config *v1alpha1.EdgeController, factory k8sinformer. uc.queryNodeChan = make(chan model.Message, config.Buffer.QueryNode) uc.updateNodeChan = make(chan model.Message, config.Buffer.UpdateNode) uc.patchPodChan = make(chan model.Message, config.Buffer.PatchPod) + uc.createPodChan = make(chan model.Message, config.Buffer.CreatePod) uc.podDeleteChan = make(chan model.Message, config.Buffer.DeletePod) uc.createLeaseChan = make(chan model.Message, config.Buffer.CreateLease) uc.queryLeaseChan = make(chan model.Message, config.Buffer.QueryLease) |
