Skip to content

Commit cef45ff

Browse files
committed
improve deploy job when strategy is merge
Signed-off-by: Patrick Zhao <zhaoyu@koderover.com>
1 parent 316664e commit cef45ff

10 files changed

Lines changed: 344 additions & 38 deletions

File tree

pkg/microservice/aslan/core/common/repository/models/wokflow_task_v4.go

Lines changed: 3 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -223,6 +223,7 @@ type JobTaskDeploySpec struct {
223223
KeyVals []*ServiceKeyVal `bson:"key_vals" json:"key_vals" yaml:"key_vals"` // deprecated since 1.18.0
224224
VariableKVs []*commontypes.RenderVariableKV `bson:"variable_kvs" json:"variable_kvs" yaml:"variable_kvs"` // new since 1.18.0, only used for k8s
225225
UpdateConfig bool `bson:"update_config" json:"update_config" yaml:"update_config"`
226+
IsImportToDeploy bool `bson:"is_import_to_deploy" json:"is_import_to_deploy" yaml:"is_import_to_deploy"`
226227
YamlContent string `bson:"yaml_content" json:"yaml_content" yaml:"yaml_content"`
227228
ServiceAndImages []*DeployServiceModule `bson:"service_and_images" json:"service_and_images" yaml:"service_and_images"`
228229
ServiceType string `bson:"service_type" json:"service_type" yaml:"service_type"`
@@ -234,7 +235,7 @@ type JobTaskDeploySpec struct {
234235
ReplaceResources []Resource `bson:"replace_resources" json:"replace_resources" yaml:"replace_resources"`
235236
RelatedPodLabels []map[string]string `bson:"-" json:"-" yaml:"-"`
236237
// overrideResource is used to do a full yaml override instead of a 2-way merge patching for all the resources
237-
OverrideResource bool `bson:"override_resource" json:"override_resource" yaml:"override_resource"`
238+
OverrideResource bool `bson:"override_resource" json:"override_resource" yaml:"override_resource"`
238239
// for compatibility
239240
ServiceModule string `bson:"service_module" json:"service_module" yaml:"-"`
240241
Image string `bson:"image" json:"image" yaml:"-"`
@@ -266,7 +267,7 @@ type JobTaskDeployRevertSpec struct {
266267
OverrideKVs string `bson:"override_kvs" json:"override_kvs" yaml:"override_kvs"`
267268
Revision int64 `bson:"revision" json:"revision" yaml:"revision"`
268269
RevisionCreateTime int64 `bson:"revision_create_time" json:"revision_create_time" yaml:"revision_create_time"`
269-
OverrideResource bool `bson:"override_resource" json:"override_resource" yaml:"override_resource"`
270+
OverrideResource bool `bson:"override_resource" json:"override_resource" yaml:"override_resource"`
270271
}
271272

272273
type DeployServiceModule struct {

pkg/microservice/aslan/core/common/repository/models/workflow_v4.go

Lines changed: 7 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -544,12 +544,13 @@ type DeployHelmChart struct {
544544
}
545545

546546
type DeployBasicInfo struct {
547-
ServiceName string `bson:"service_name" yaml:"service_name" json:"service_name"`
548-
Modules []*DeployModuleInfo `bson:"modules" yaml:"modules" json:"modules"`
549-
Deployed bool `bson:"deployed" yaml:"deployed" json:"deployed"`
550-
AutoSync bool `bson:"-" yaml:"auto_sync" json:"auto_sync"`
551-
UpdateConfig bool `bson:"update_config" yaml:"update_config" json:"update_config"`
552-
Updatable bool `bson:"-" yaml:"updatable" json:"updatable"`
547+
ServiceName string `bson:"service_name" yaml:"service_name" json:"service_name"`
548+
DeployStrategy string `bson:"deploy_strategy" yaml:"deploy_strategy" json:"deploy_strategy"`
549+
Modules []*DeployModuleInfo `bson:"modules" yaml:"modules" json:"modules"`
550+
Deployed bool `bson:"deployed" yaml:"deployed" json:"deployed"`
551+
AutoSync bool `bson:"-" yaml:"auto_sync" json:"auto_sync"`
552+
UpdateConfig bool `bson:"update_config" yaml:"update_config" json:"update_config"`
553+
Updatable bool `bson:"-" yaml:"updatable" json:"updatable"`
553554
}
554555

555556
type DeployOptionInfo struct {

pkg/microservice/aslan/core/common/service/kube/render.go

Lines changed: 275 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -32,7 +32,7 @@ import (
3232
yamlutil "k8s.io/apimachinery/pkg/util/yaml"
3333
"k8s.io/cli-runtime/pkg/printers"
3434
"k8s.io/helm/pkg/releaseutil"
35-
yaml "sigs.k8s.io/yaml/goyaml.v3"
35+
"sigs.k8s.io/yaml"
3636

3737
"github.com/koderover/zadig/v2/pkg/microservice/aslan/config"
3838
"github.com/koderover/zadig/v2/pkg/microservice/aslan/core/common/repository/models"
@@ -58,6 +58,7 @@ type GeneSvcYamlOption struct {
5858
EnvName string
5959
ServiceName string
6060
UpdateServiceRevision bool
61+
IsImportToDeploy bool
6162
VariableYaml string
6263
VariableKVs []*commontypes.RenderVariableKV
6364
ReplicaOverrides []*commonmodels.WorkLoad
@@ -400,20 +401,284 @@ func FetchCurrentAppliedYaml(option *GeneSvcYamlOption) (string, int, error) {
400401
return "", 0, errors.Wrapf(err, "failed to find service %s with revision %d", option.ServiceName, curProductSvc.Revision)
401402
}
402403

403-
fullRenderedYaml, err := RenderServiceYaml(prodSvcTemplate.Yaml, option.ProductName, option.ServiceName, curProductSvc.GetServiceRender())
404+
if option.IsImportToDeploy {
405+
importedAllManifests, _, err := FetchImportedAllManifests(productInfo, prodSvcTemplate, curProductSvc.GetServiceRender())
406+
if err != nil {
407+
return "", 0, err
408+
}
409+
return importedAllManifests, 0, err
410+
} else {
411+
fullRenderedYaml, err := RenderServiceYaml(prodSvcTemplate.Yaml, option.ProductName, option.ServiceName, curProductSvc.GetServiceRender())
412+
if err != nil {
413+
return "", 0, err
414+
}
415+
fullRenderedYaml = ParseSysKeys(productInfo.Namespace, productInfo.EnvName, option.ProductName, option.ServiceName, fullRenderedYaml)
416+
mergedContainers := mergeContainers(prodSvcTemplate.Containers, curProductSvc.Containers)
417+
fullRenderedYaml, _, err = ReplaceWorkloadImages(fullRenderedYaml, mergedContainers)
418+
if err != nil {
419+
return "", 0, err
420+
}
421+
422+
replicaOverrides := resolveReplicaOverrides(curProductSvc.WorkLoads, option.ReplicaOverrides, option.IgnoreCurrentReplicaOverrides)
423+
fullRenderedYaml, err = ApplyReplicaOverrides(fullRenderedYaml, replicaOverrides)
424+
return fullRenderedYaml, 0, err
425+
}
426+
}
427+
428+
func FetchImportedAllManifests(envInfo *models.Product, serviceTmp *models.Service, svcRender *template.ServiceRender) (string, []*WorkloadResource, error) {
429+
fullRenderedYaml, err := RenderServiceYaml(serviceTmp.Yaml, envInfo.ProductName, serviceTmp.ServiceName, svcRender)
404430
if err != nil {
405-
return "", 0, err
431+
return "", nil, err
406432
}
407-
fullRenderedYaml = ParseSysKeys(productInfo.Namespace, productInfo.EnvName, option.ProductName, option.ServiceName, fullRenderedYaml)
408-
mergedContainers := mergeContainers(prodSvcTemplate.Containers, curProductSvc.Containers)
409-
fullRenderedYaml, _, err = ReplaceWorkloadImages(fullRenderedYaml, mergedContainers)
433+
fullRenderedYaml = ParseSysKeys(envInfo.Namespace, envInfo.EnvName, envInfo.ProductName, serviceTmp.ServiceName, fullRenderedYaml)
434+
435+
manifests := releaseutil.SplitManifests(fullRenderedYaml)
436+
437+
kubeClient, err := clientmanager.NewKubeClientManager().GetControllerRuntimeClient(envInfo.ClusterID)
410438
if err != nil {
411-
return "", 0, err
439+
log.Errorf("cluster is not connected [%s]", envInfo.ClusterID)
440+
return "", nil, errors.Wrapf(err, "cluster is not connected [%s]", envInfo.ClusterID)
412441
}
413442

414-
replicaOverrides := resolveReplicaOverrides(curProductSvc.WorkLoads, option.ReplicaOverrides, option.IgnoreCurrentReplicaOverrides)
415-
fullRenderedYaml, err = ApplyReplicaOverrides(fullRenderedYaml, replicaOverrides)
416-
return fullRenderedYaml, 0, err
443+
clientset, err := clientmanager.NewKubeClientManager().GetKubernetesClientSet(envInfo.ClusterID)
444+
if err != nil {
445+
log.Errorf("get client set error: %v", err)
446+
return "", nil, err
447+
}
448+
versionInfo, err := clientset.Discovery().ServerVersion()
449+
if err != nil {
450+
log.Errorf("get server version error: %v", err)
451+
return "", nil, err
452+
}
453+
manifestArr := make([]string, 0)
454+
workloadRes := make([]*WorkloadResource, 0)
455+
456+
for _, item := range manifests {
457+
u, err := serializer.NewDecoder().YamlToUnstructured([]byte(item))
458+
if err != nil {
459+
return "", nil, errors.Wrapf(err, "failed to decode yaml %s", item)
460+
}
461+
kind := u.GetKind()
462+
name := u.GetName()
463+
464+
var bs []byte
465+
var exist bool
466+
isWorkload := false
467+
468+
switch kind {
469+
case setting.Deployment:
470+
bs, exist, err = getter.GetDeploymentYamlFormat(envInfo.Namespace, name, kubeClient)
471+
isWorkload = true
472+
case setting.StatefulSet:
473+
bs, exist, err = getter.GetStatefulSetYamlFormat(envInfo.Namespace, name, kubeClient)
474+
isWorkload = true
475+
case setting.CronJob:
476+
bs, exist, err = getter.GetCronJobYamlFormat(envInfo.Namespace, name, kubeClient, kubeclient.VersionLessThan121(versionInfo))
477+
isWorkload = true
478+
case setting.Job:
479+
bs, exist, err = getter.GetJobYaml(envInfo.Namespace, name, kubeClient)
480+
isWorkload = true
481+
case setting.Service:
482+
bs, exist, err = getter.GetServiceYamlFormat(envInfo.Namespace, name, kubeClient)
483+
case setting.ConfigMap:
484+
bs, exist, err = getter.GetConfigMapYamlFormat(envInfo.Namespace, name, kubeClient)
485+
case setting.Secret:
486+
bs, exist, err = getter.GetSecretYamlFormat(envInfo.Namespace, name, kubeClient)
487+
case setting.Ingress:
488+
bs, exist, err = getter.GetIngressYamlFormat(envInfo.Namespace, name, kubeClient)
489+
case setting.PersistentVolumeClaim:
490+
bs, exist, err = getter.GetPVCYamlFormat(envInfo.Namespace, name, kubeClient)
491+
default:
492+
log.Warnf("unsupported resource kind %s/%s, skipping", kind, name)
493+
continue
494+
}
495+
496+
if err != nil {
497+
return "", nil, errors.Wrapf(err, "failed to get %s %s", kind, name)
498+
}
499+
if !exist {
500+
return "", nil, errors.Errorf("%s %s not found", kind, name)
501+
}
502+
if isWorkload {
503+
workloadRes = append(workloadRes, &WorkloadResource{Type: kind, Name: name})
504+
}
505+
506+
cleaned, err := cleanupClusterResource(bs)
507+
if err != nil {
508+
return "", nil, errors.Wrapf(err, "failed to cleanup %s %s", kind, name)
509+
}
510+
manifestArr = append(manifestArr, cleaned)
511+
}
512+
513+
return util.JoinYamls(manifestArr), workloadRes, nil
514+
}
515+
516+
var serverGeneratedAnnotations = []string{
517+
"kubectl.kubernetes.io/last-applied-configuration",
518+
"deployment.kubernetes.io/revision",
519+
"deprecated.daemonset.template.generation",
520+
"pv.kubernetes.io/bind-completed",
521+
"pv.kubernetes.io/bound-by-controller",
522+
}
523+
524+
func cleanupClusterResource(yamlData []byte) (string, error) {
525+
obj, err := serializer.NewDecoder().YamlToUnstructured(yamlData)
526+
if err != nil {
527+
return "", err
528+
}
529+
530+
delete(obj.Object, "status")
531+
532+
if metadata, ok := obj.Object["metadata"].(map[string]interface{}); ok {
533+
delete(metadata, "managedFields")
534+
delete(metadata, "resourceVersion")
535+
delete(metadata, "uid")
536+
delete(metadata, "selfLink")
537+
delete(metadata, "creationTimestamp")
538+
delete(metadata, "generation")
539+
}
540+
541+
annotations := obj.GetAnnotations()
542+
if annotations != nil {
543+
for _, key := range serverGeneratedAnnotations {
544+
delete(annotations, key)
545+
}
546+
if len(annotations) == 0 {
547+
obj.SetAnnotations(nil)
548+
} else {
549+
obj.SetAnnotations(annotations)
550+
}
551+
}
552+
553+
cleanupDefaultedSpec(obj.Object, obj.GetKind())
554+
555+
resp, err := yaml.Marshal(obj.Object)
556+
if err != nil {
557+
return "", err
558+
}
559+
return string(resp), nil
560+
}
561+
562+
func cleanupDefaultedSpec(objMap map[string]interface{}, kind string) {
563+
cleanupPodTemplate(objMap, "spec", "template")
564+
cleanupPodTemplate(objMap, "spec", "jobTemplate", "spec", "template")
565+
566+
if jobTmpl := nestedMap(objMap, "spec", "jobTemplate"); jobTmpl != nil {
567+
if metadata, ok := jobTmpl["metadata"].(map[string]interface{}); ok {
568+
delete(metadata, "creationTimestamp")
569+
}
570+
}
571+
572+
switch kind {
573+
case setting.Service:
574+
cleanupServiceDefaults(objMap)
575+
case setting.PersistentVolumeClaim:
576+
if spec := nestedMap(objMap, "spec"); spec != nil {
577+
deleteIfDefaultStr(spec, "volumeMode", "Filesystem")
578+
}
579+
case setting.Secret:
580+
deleteIfDefaultStr(objMap, "type", "Opaque")
581+
}
582+
}
583+
584+
func cleanupPodTemplate(objMap map[string]interface{}, path ...string) {
585+
tmpl := nestedMap(objMap, path...)
586+
if tmpl == nil {
587+
return
588+
}
589+
590+
if metadata, ok := tmpl["metadata"].(map[string]interface{}); ok {
591+
delete(metadata, "creationTimestamp")
592+
}
593+
594+
podSpec, ok := tmpl["spec"].(map[string]interface{})
595+
if !ok {
596+
return
597+
}
598+
599+
deleteIfDefaultStr(podSpec, "dnsPolicy", "ClusterFirst")
600+
deleteIfDefaultStr(podSpec, "schedulerName", "default-scheduler")
601+
deleteIfDefaultBool(podSpec, "enableServiceLinks", true)
602+
deleteEmptyMap(podSpec, "securityContext")
603+
delete(podSpec, "serviceAccount")
604+
605+
for _, key := range []string{"containers", "initContainers"} {
606+
containers, ok := podSpec[key].([]interface{})
607+
if !ok {
608+
continue
609+
}
610+
for _, c := range containers {
611+
container, ok := c.(map[string]interface{})
612+
if !ok {
613+
continue
614+
}
615+
deleteIfDefaultStr(container, "terminationMessagePath", "/dev/termination-log")
616+
deleteIfDefaultStr(container, "terminationMessagePolicy", "File")
617+
deleteEmptyMap(container, "resources")
618+
if ports, ok := container["ports"].([]interface{}); ok {
619+
for _, p := range ports {
620+
if port, ok := p.(map[string]interface{}); ok {
621+
deleteIfDefaultStr(port, "protocol", "TCP")
622+
}
623+
}
624+
}
625+
}
626+
}
627+
}
628+
629+
func cleanupServiceDefaults(objMap map[string]interface{}) {
630+
spec := nestedMap(objMap, "spec")
631+
if spec == nil {
632+
return
633+
}
634+
635+
if clusterIP, ok := spec["clusterIP"].(string); ok && clusterIP != "None" {
636+
delete(spec, "clusterIP")
637+
delete(spec, "clusterIPs")
638+
}
639+
640+
deleteIfDefaultStr(spec, "sessionAffinity", "None")
641+
deleteIfDefaultStr(spec, "internalTrafficPolicy", "Cluster")
642+
delete(spec, "ipFamilies")
643+
delete(spec, "ipFamilyPolicy")
644+
645+
if ports, ok := spec["ports"].([]interface{}); ok {
646+
for _, p := range ports {
647+
if port, ok := p.(map[string]interface{}); ok {
648+
deleteIfDefaultStr(port, "protocol", "TCP")
649+
}
650+
}
651+
}
652+
}
653+
654+
func nestedMap(obj map[string]interface{}, keys ...string) map[string]interface{} {
655+
cur := obj
656+
for _, k := range keys {
657+
next, ok := cur[k].(map[string]interface{})
658+
if !ok {
659+
return nil
660+
}
661+
cur = next
662+
}
663+
return cur
664+
}
665+
666+
func deleteIfDefaultStr(m map[string]interface{}, key, defaultVal string) {
667+
if v, ok := m[key].(string); ok && v == defaultVal {
668+
delete(m, key)
669+
}
670+
}
671+
672+
func deleteIfDefaultBool(m map[string]interface{}, key string, defaultVal bool) {
673+
if v, ok := m[key].(bool); ok && v == defaultVal {
674+
delete(m, key)
675+
}
676+
}
677+
678+
func deleteEmptyMap(m map[string]interface{}, key string) {
679+
if v, ok := m[key].(map[string]interface{}); ok && len(v) == 0 {
680+
delete(m, key)
681+
}
417682
}
418683

419684
func FetchImportedManifests(option *GeneSvcYamlOption, productInfo *models.Product, serviceTmp *models.Service, svcRender *template.ServiceRender) (string, []*WorkloadResource, error) {

0 commit comments

Comments
 (0)