From d81a259096d4a0be78e1e8168a89464cd4fd17a0 Mon Sep 17 00:00:00 2001 From: huanghong Date: Mon, 1 Dec 2025 17:20:27 +0800 Subject: [PATCH] =?UTF-8?q?feat(dtc-install):=20=E5=A2=9E=E5=8A=A0?= =?UTF-8?q?=E5=AD=90chart=E5=AE=89=E8=A3=85=E9=A1=BA=E5=BA=8F?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- cmd/helm/install.go | 4 + pkg/action/action.go | 128 +++++++ pkg/action/install.go | 587 +++++++++++++++++++++++++++++ pkg/action/install_sequence.go | 152 ++++++++ pkg/chart/chart.go | 21 ++ pkg/chart/metadata.go | 7 + pkg/kube/client.go | 5 + pkg/releaseutil/chart_sorter.go | 94 +++++ pkg/releaseutil/manifest_sorter.go | 39 ++ 9 files changed, 1037 insertions(+) create mode 100644 pkg/action/install_sequence.go create mode 100644 pkg/releaseutil/chart_sorter.go diff --git a/cmd/helm/install.go b/cmd/helm/install.go index 6ffc968ce..e883b3d74 100644 --- a/cmd/helm/install.go +++ b/cmd/helm/install.go @@ -284,6 +284,10 @@ func runInstall(args []string, client *action.Install, valueOpts *values.Options cancel() }() + if chartRequested.IsSequenceInstall() { + return client.RunWithCustomContext(ctx, chartRequested, vals) + } + return client.RunWithContext(ctx, chartRequested, vals) } diff --git a/pkg/action/action.go b/pkg/action/action.go index 016aec3f6..6f336248e 100644 --- a/pkg/action/action.go +++ b/pkg/action/action.go @@ -230,6 +230,134 @@ func (cfg *Configuration) renderResources(ch *chart.Chart, values chartutil.Valu return hs, b, notes, nil } +func (cfg *Configuration) renderResourcesWithChartOrder(ch *chart.Chart, values chartutil.Values, releaseName, outputDir string, subNotes, useReleaseName, includeCrds bool, pr postrender.PostRenderer, dryRun, enableDNS bool) ([]*release.Hook, *bytes.Buffer, string, []releaseutil.Manifest, error) { + hs := []*release.Hook{} + b := bytes.NewBuffer(nil) + + caps, err := cfg.getCapabilities() + if err != nil { + return hs, b, "", nil, err + } + + if ch.Metadata.KubeVersion != "" { + if !chartutil.IsCompatibleRange(ch.Metadata.KubeVersion, caps.KubeVersion.String()) { + return hs, b, "", nil, errors.Errorf("chart requires kubeVersion: %s which is incompatible with Kubernetes %s", ch.Metadata.KubeVersion, caps.KubeVersion.String()) + } + } + + var files map[string]string + var err2 error + + // A `helm template` or `helm install --dry-run` should not talk to the remote cluster. + // It will break in interesting and exotic ways because other data (e.g. discovery) + // is mocked. It is not up to the template author to decide when the user wants to + // connect to the cluster. So when the user says to dry run, respect the user's + // wishes and do not connect to the cluster. + if !dryRun && cfg.RESTClientGetter != nil { + restConfig, err := cfg.RESTClientGetter.ToRESTConfig() + if err != nil { + return hs, b, "", nil, err + } + e := engine.New(restConfig) + e.EnableDNS = enableDNS + files, err2 = e.Render(ch, values) + } else { + var e engine.Engine + e.EnableDNS = enableDNS + files, err2 = e.Render(ch, values) + } + + if err2 != nil { + return hs, b, "", nil, err2 + } + + // NOTES.txt gets rendered like all the other files, but because it's not a hook nor a resource, + // pull it out of here into a separate file so that we can actually use the output of the rendered + // text file. We have to spin through this map because the file contains path information, so we + // look for terminating NOTES.txt. We also remove it from the files so that we don't have to skip + // it in the sortHooks. + var notesBuffer bytes.Buffer + for k, v := range files { + if strings.HasSuffix(k, notesFileSuffix) { + if subNotes || (k == path.Join(ch.Name(), "templates", notesFileSuffix)) { + // If buffer contains data, add newline before adding more + if notesBuffer.Len() > 0 { + notesBuffer.WriteString("\n") + } + notesBuffer.WriteString(v) + } + delete(files, k) + } + } + notes := notesBuffer.String() + + // Sort hooks, manifests, and partials. Only hooks and manifests are returned, + // as partials are not used after renderer.Render. Empty manifests are also + // removed here. + chartInstallOrder := ch.ChartInstallOrder() + hs, manifests, err := releaseutil.SortManifestsByChart(files, caps.APIVersions, chartInstallOrder) + if err != nil { + // By catching parse errors here, we can prevent bogus releases from going + // to Kubernetes. + // + // We return the files as a big blob of data to help the user debug parser + // errors. + for name, content := range files { + if strings.TrimSpace(content) == "" { + continue + } + fmt.Fprintf(b, "---\n# Source: %s\n%s\n", name, content) + } + return hs, b, "", nil, err + } + + // Aggregate all valid manifests into one big doc. + fileWritten := make(map[string]bool) + + if includeCrds { + for _, crd := range ch.CRDObjects() { + if outputDir == "" { + fmt.Fprintf(b, "---\n# Source: %s\n%s\n", crd.Name, string(crd.File.Data[:])) + } else { + err = writeToFile(outputDir, crd.Filename, string(crd.File.Data[:]), fileWritten[crd.Name]) + if err != nil { + return hs, b, "", nil, err + } + fileWritten[crd.Name] = true + } + } + } + + for _, m := range manifests { + if outputDir == "" { + fmt.Fprintf(b, "---\n# Source: %s\n%s\n", m.Name, m.Content) + } else { + newDir := outputDir + if useReleaseName { + newDir = filepath.Join(outputDir, releaseName) + } + // NOTE: We do not have to worry about the post-renderer because + // output dir is only used by `helm template`. In the next major + // release, we should move this logic to template only as it is not + // used by install or upgrade + err = writeToFile(newDir, m.Name, m.Content, fileWritten[m.Name]) + if err != nil { + return hs, b, "", nil, err + } + fileWritten[m.Name] = true + } + } + + if pr != nil { + b, err = pr.Run(b) + if err != nil { + return hs, b, notes, nil, errors.Wrap(err, "error while running post render on files") + } + } + + return hs, b, notes, manifests, nil +} + // RESTClientGetter gets the rest client type RESTClientGetter interface { ToRESTConfig() (*rest.Config, error) diff --git a/pkg/action/install.go b/pkg/action/install.go index d19972e1b..1bca3c7c8 100644 --- a/pkg/action/install.go +++ b/pkg/action/install.go @@ -23,6 +23,7 @@ import ( "io" "net/url" "os" + "os/exec" "path" "path/filepath" "strings" @@ -127,6 +128,13 @@ type ChartPathOptions struct { registryClient *registry.Client } +type InstallItem struct { + ChartName string + WaitFor string + Manifests []releaseutil.Manifest + Resources kube.ResourceList +} + // NewInstall creates a new Install object with the given configuration. func NewInstall(cfg *Configuration) *Install { in := &Install{ @@ -167,6 +175,7 @@ func (i *Install) installCRDs(crds []chart.CRD) error { } return errors.Wrapf(err, "failed to install CRD %s", obj.Name) } + i.cfg.Log("Installing CRD: %s", obj.Filename) totalItems = append(totalItems, res...) } if len(totalItems) > 0 { @@ -387,6 +396,478 @@ func (i *Install) RunWithContext(ctx context.Context, chrt *chart.Chart, vals ma } } +//func (i *Install) RunWithCustomContext(ctx context.Context, chrt *chart.Chart, vals map[string]interface{}) (*release.Release, error) { +// // Check reachability of cluster unless in client-only mode (e.g. `helm template` without `--validate`) +// if !i.ClientOnly { +// if err := i.cfg.KubeClient.IsReachable(); err != nil { +// return nil, err +// } +// } +// +// if err := i.availableName(); err != nil { +// return nil, err +// } +// +// if err := chartutil.ProcessDependencies(chrt, vals); err != nil { +// return nil, err +// } +// +// // Pre-install anything in the crd/ directory. We do this before Helm +// // contacts the upstream server and builds the capabilities object. +// // 先install 所有的crd +// if crds := chrt.CRDObjects(); !i.ClientOnly && !i.SkipCRDs && len(crds) > 0 { +// // On dry run, bail here +// if i.DryRun { +// i.cfg.Log("WARNING: This chart or one of its subcharts contains CRDs. Rendering may fail or contain inaccuracies.") +// } else if err := i.installCRDs(crds); err != nil { +// return nil, err +// } +// } +// +// if i.ClientOnly { +// // Add mock objects in here so it doesn't use Kube API server +// // NOTE(bacongobbler): used for `helm template` +// i.cfg.Capabilities = chartutil.DefaultCapabilities.Copy() +// if i.KubeVersion != nil { +// i.cfg.Capabilities.KubeVersion = *i.KubeVersion +// } +// i.cfg.Capabilities.APIVersions = append(i.cfg.Capabilities.APIVersions, i.APIVersions...) +// i.cfg.KubeClient = &kubefake.PrintingKubeClient{Out: io.Discard} +// +// mem := driver.NewMemory() +// mem.SetNamespace(i.Namespace) +// i.cfg.Releases = storage.Init(mem) +// } else if !i.ClientOnly && len(i.APIVersions) > 0 { +// i.cfg.Log("API Version list given outside of client only mode, this list will be ignored") +// } +// +// // Make sure if Atomic is set, that wait is set as well. This makes it so +// // the user doesn't have to specify both +// i.Wait = i.Wait || i.Atomic +// +// caps, err := i.cfg.getCapabilities() +// if err != nil { +// return nil, err +// } +// +// // special case for helm template --is-upgrade +// isUpgrade := i.IsUpgrade && i.DryRun +// options := chartutil.ReleaseOptions{ +// Name: i.ReleaseName, +// Namespace: i.Namespace, +// Revision: 1, +// IsInstall: !isUpgrade, +// IsUpgrade: isUpgrade, +// } +// // 汇总的values文件 +// valuesToRender, err := chartutil.ToRenderValues(chrt, vals, options, caps) +// if err != nil { +// return nil, err +// } +// +// rel := i.createRelease(chrt, vals) +// +// var InstallDeps []chart.InstallItem +// if chrt.IsSequenceInstall() { +// InstallDeps, err = chrt.GetInstallItems(chrt.Metadata.Install) +// if err != nil { +// return nil, err +// } +// } +// // 按照顺序执行安装 +// allmanifestDoc := new(bytes.Buffer) +// var allhooks []*release.Hook +// +// for _, seq := range InstallDeps { +// var manifestDoc *bytes.Buffer +// var hooks []*release.Hook +// +// hooks, manifestDoc, _, err = i.cfg.renderResources(seq.Chart, valuesToRender, i.ReleaseName, i.OutputDir, i.SubNotes, i.UseReleaseName, i.IncludeCRDs, i.PostRenderer, i.DryRun, i.EnableDNS) +// // Check error from render +// if err != nil { +// rel.SetStatus(release.StatusFailed, fmt.Sprintf("failed to render resource: %s", err.Error())) +// // Return a release with partial data so that the client can show debugging information. +// return rel, err +// } +// +// allhooks = append(allhooks, hooks...) +// if manifestDoc != nil { +// allmanifestDoc.Write(manifestDoc.Bytes()) +// } +// +// // Mark this release as in-progress +// rel.SetStatus(release.StatusPendingInstall, "Initial install underway") +// +// var toBeAdopted kube.ResourceList +// resources, err := i.cfg.KubeClient.Build(bytes.NewBufferString(manifestDoc.String()), !i.DisableOpenAPIValidation) +// if err != nil { +// return nil, errors.Wrap(err, "unable to build kubernetes objects from release manifest") +// } +// +// // It is safe to use "force" here because these are resources currently rendered by the chart. +// err = resources.Visit(setMetadataVisitor(rel.Name, rel.Namespace, true)) +// if err != nil { +// return nil, err +// } +// +// if !i.ClientOnly && !isUpgrade && len(resources) > 0 { +// toBeAdopted, err = existingResourceConflict(resources, rel.Name, rel.Namespace) +// if err != nil { +// return nil, errors.Wrap(err, "rendered manifests contain a resource that already exists. Unable to continue with install") +// } +// } +// +// if i.CreateNamespace { +// ns := &v1.Namespace{ +// TypeMeta: metav1.TypeMeta{ +// APIVersion: "v1", +// Kind: "Namespace", +// }, +// ObjectMeta: metav1.ObjectMeta{ +// Name: i.Namespace, +// Labels: map[string]string{ +// "name": i.Namespace, +// }, +// }, +// } +// buf, err := yaml.Marshal(ns) +// if err != nil { +// return nil, err +// } +// resourceList, err := i.cfg.KubeClient.Build(bytes.NewBuffer(buf), true) +// if err != nil { +// return nil, err +// } +// if _, err := i.cfg.KubeClient.Create(resourceList); err != nil && !apierrors.IsAlreadyExists(err) { +// return nil, err +// } +// } +// +// if len(toBeAdopted) == 0 && len(resources) > 0 { +// if _, err = i.cfg.KubeClient.Create(resources); err != nil { +// return nil, err +// } +// } else if len(resources) > 0 { +// if _, err = i.cfg.KubeClient.Update(toBeAdopted, resources, i.Force); err != nil { +// return nil, err +// } +// } +// +// if seq.WaitFor != "" { +// i.cfg.Log("Waiting for resources: %s", seq.WaitFor) +// +// cmd := exec.Command("bash", "-c", seq.WaitFor) +// stdoutStderr, err := cmd.CombinedOutput() +// +// if err != nil { +// // kubectl wait 失败(超时或资源不存在等) +// i.cfg.Log("Wait command failed: %s, output: %s", err.Error(), string(stdoutStderr)) +// rel.SetStatus(release.StatusFailed, fmt.Sprintf("wait command failed: %s", string(stdoutStderr))) +// return rel, fmt.Errorf("wait for resource failed: %s", string(stdoutStderr)) +// } +// +// // 成功,继续安装后续 chart +// i.cfg.Log("Wait completed: %s", string(stdoutStderr)) +// } +// +// } +// +// // Even for errors, attach this if available +// if allmanifestDoc != nil { +// rel.Manifest = allmanifestDoc.String() +// } +// rel.Hooks = allhooks +// +// if err := i.cfg.Releases.Create(rel); err != nil { +// // We could try to recover gracefully here, but since nothing has been installed +// // yet, this is probably safer than trying to continue when we know storage is +// // not working. +// return rel, err +// } +// return rel, nil +// +// //var manifestDoc *bytes.Buffer +// //rel.Hooks, manifestDoc, rel.Info.Notes, err = i.cfg.renderResources(chrt, valuesToRender, i.ReleaseName, i.OutputDir, i.SubNotes, i.UseReleaseName, i.IncludeCRDs, i.PostRenderer, i.DryRun, i.EnableDNS) +// // Even for errors, attach this if available +// //if manifestDoc != nil { +// // rel.Manifest = manifestDoc.String() +// //} +// // Check error from render +// //if err != nil { +// // rel.SetStatus(release.StatusFailed, fmt.Sprintf("failed to render resource: %s", err.Error())) +// // Return a release with partial data so that the client can show debugging information. +// //return rel, err +// //} +// +// //// Mark this release as in-progress +// //rel.SetStatus(release.StatusPendingInstall, "Initial install underway") +// +// //var toBeAdopted kube.ResourceList +// //resources, err := i.cfg.KubeClient.Build(bytes.NewBufferString(rel.Manifest), !i.DisableOpenAPIValidation) +// //if err != nil { +// // return nil, errors.Wrap(err, "unable to build kubernetes objects from release manifest") +// //} +// // +// //// It is safe to use "force" here because these are resources currently rendered by the chart. +// //err = resources.Visit(setMetadataVisitor(rel.Name, rel.Namespace, true)) +// //if err != nil { +// // return nil, err +// //} +// +// // Install requires an extra validation step of checking that resources +// // don't already exist before we actually create resources. If we continue +// // forward and create the release object with resources that already exist, +// // we'll end up in a state where we will delete those resources upon +// // deleting the release because the manifest will be pointing at that +// // resource +// //if !i.ClientOnly && !isUpgrade && len(resources) > 0 { +// // toBeAdopted, err = existingResourceConflict(resources, rel.Name, rel.Namespace) +// // if err != nil { +// // return nil, errors.Wrap(err, "rendered manifests contain a resource that already exists. Unable to continue with install") +// // } +// //} +// +// // Bail out here if it is a dry run +// //######## dryrun失效 +// //if i.DryRun { +// // rel.Info.Description = "Dry run complete" +// // return rel, nil +// //} +// // +// //if i.CreateNamespace { +// // ns := &v1.Namespace{ +// // TypeMeta: metav1.TypeMeta{ +// // APIVersion: "v1", +// // Kind: "Namespace", +// // }, +// // ObjectMeta: metav1.ObjectMeta{ +// // Name: i.Namespace, +// // Labels: map[string]string{ +// // "name": i.Namespace, +// // }, +// // }, +// // } +// // buf, err := yaml.Marshal(ns) +// // if err != nil { +// // return nil, err +// // } +// // resourceList, err := i.cfg.KubeClient.Build(bytes.NewBuffer(buf), true) +// // if err != nil { +// // return nil, err +// // } +// // if _, err := i.cfg.KubeClient.Create(resourceList); err != nil && !apierrors.IsAlreadyExists(err) { +// // return nil, err +// // } +// //} +// +// // If Replace is true, we need to supercede the last release. +// // Release 无效 +// //if i.Replace { +// // if err := i.replaceRelease(rel); err != nil { +// // return nil, err +// // } +// //} +// +// //// Store the release in history before continuing (new in Helm 3). We always know +// //// that this is a create operation. +// //if err := i.cfg.Releases.Create(rel); err != nil { +// // // We could try to recover gracefully here, but since nothing has been installed +// // // yet, this is probably safer than trying to continue when we know storage is +// // // not working. +// // return rel, err +// //} +// //rChan := make(chan resultMessage) +// //ctxChan := make(chan resultMessage) +// //doneChan := make(chan struct{}) +// //defer close(doneChan) +// //go i.performInstall(rChan, rel, toBeAdopted, resources) +// //go i.handleContext(ctx, ctxChan, doneChan, rel) +// //select { +// //case result := <-rChan: +// // return result.r, result.e +// //case result := <-ctxChan: +// // return result.r, result.e +// //} +//} + +func (i *Install) RunWithCustomContext(ctx context.Context, chrt *chart.Chart, vals map[string]interface{}) (*release.Release, error) { + // Check reachability of cluster unless in client-only mode (e.g. `helm template` without `--validate`) + if !i.ClientOnly { + if err := i.cfg.KubeClient.IsReachable(); err != nil { + return nil, err + } + } + + if err := i.availableName(); err != nil { + return nil, err + } + + if err := chartutil.ProcessDependencies(chrt, vals); err != nil { + return nil, err + } + + // Pre-install anything in the crd/ directory. We do this before Helm + // contacts the upstream server and builds the capabilities object. + if crds := chrt.CRDObjects(); !i.ClientOnly && !i.SkipCRDs && len(crds) > 0 { + // On dry run, bail here + if i.DryRun { + i.cfg.Log("WARNING: This chart or one of its subcharts contains CRDs. Rendering may fail or contain inaccuracies.") + } else if err := i.installCRDs(crds); err != nil { + return nil, err + } + } + + if i.ClientOnly { + // Add mock objects in here so it doesn't use Kube API server + // NOTE(bacongobbler): used for `helm template` + i.cfg.Capabilities = chartutil.DefaultCapabilities.Copy() + if i.KubeVersion != nil { + i.cfg.Capabilities.KubeVersion = *i.KubeVersion + } + i.cfg.Capabilities.APIVersions = append(i.cfg.Capabilities.APIVersions, i.APIVersions...) + i.cfg.KubeClient = &kubefake.PrintingKubeClient{Out: io.Discard} + + mem := driver.NewMemory() + mem.SetNamespace(i.Namespace) + i.cfg.Releases = storage.Init(mem) + } else if !i.ClientOnly && len(i.APIVersions) > 0 { + i.cfg.Log("API Version list given outside of client only mode, this list will be ignored") + } + + // Make sure if Atomic is set, that wait is set as well. This makes it so + // the user doesn't have to specify both + i.Wait = i.Wait || i.Atomic + + caps, err := i.cfg.getCapabilities() + if err != nil { + return nil, err + } + + // special case for helm template --is-upgrade + isUpgrade := i.IsUpgrade && i.DryRun + options := chartutil.ReleaseOptions{ + Name: i.ReleaseName, + Namespace: i.Namespace, + Revision: 1, + IsInstall: !isUpgrade, + IsUpgrade: isUpgrade, + } + valuesToRender, err := chartutil.ToRenderValues(chrt, vals, options, caps) + if err != nil { + return nil, err + } + + rel := i.createRelease(chrt, vals) + + var manifestDoc *bytes.Buffer + var sortedManifest []releaseutil.Manifest + //rel.Hooks, manifestDoc, rel.Info.Notes, err = i.cfg.renderResources(chrt, valuesToRender, i.ReleaseName, i.OutputDir, i.SubNotes, i.UseReleaseName, i.IncludeCRDs, i.PostRenderer, i.DryRun, i.EnableDNS) + rel.Hooks, manifestDoc, rel.Info.Notes, sortedManifest, err = i.cfg.renderResourcesWithChartOrder(chrt, valuesToRender, i.ReleaseName, i.OutputDir, i.SubNotes, i.UseReleaseName, i.IncludeCRDs, i.PostRenderer, i.DryRun, i.EnableDNS) + // Even for errors, attach this if available + if manifestDoc != nil { + rel.Manifest = manifestDoc.String() + } + // Check error from render + if err != nil { + rel.SetStatus(release.StatusFailed, fmt.Sprintf("failed to render resource: %s", err.Error())) + // Return a release with partial data so that the client can show debugging information. + return rel, err + } + + // Mark this release as in-progress + rel.SetStatus(release.StatusPendingInstall, "Initial install underway") + if i.CreateNamespace { + ns := &v1.Namespace{ + TypeMeta: metav1.TypeMeta{ + APIVersion: "v1", + Kind: "Namespace", + }, + ObjectMeta: metav1.ObjectMeta{ + Name: i.Namespace, + Labels: map[string]string{ + "name": i.Namespace, + }, + }, + } + buf, err := yaml.Marshal(ns) + if err != nil { + return nil, err + } + resourceList, err := i.cfg.KubeClient.Build(bytes.NewBuffer(buf), true) + if err != nil { + return nil, err + } + if _, err := i.cfg.KubeClient.Create(resourceList); err != nil && !apierrors.IsAlreadyExists(err) { + return nil, err + } + } + + var toBeAdopted kube.ResourceList + resources, err := i.cfg.KubeClient.Build(bytes.NewBufferString(rel.Manifest), !i.DisableOpenAPIValidation) + if err != nil { + return nil, errors.Wrap(err, "unable to build kubernetes objects from release manifest") + } + + // It is safe to use "force" here because these are resources currently rendered by the chart. + err = resources.Visit(setMetadataVisitor(rel.Name, rel.Namespace, true)) + if err != nil { + return nil, err + } + + // Install requires an extra validation step of checking that resources + // don't already exist before we actually create resources. If we continue + // forward and create the release object with resources that already exist, + // we'll end up in a state where we will delete those resources upon + // deleting the release because the manifest will be pointing at that + // resource + if !i.ClientOnly && !isUpgrade && len(resources) > 0 { + toBeAdopted, err = existingResourceConflict(resources, rel.Name, rel.Namespace) + if err != nil { + return nil, errors.Wrap(err, "rendered manifests contain a resource that already exists. Unable to continue with install") + } + } + + // Bail out here if it is a dry run + if i.DryRun { + rel.Info.Description = "Dry run complete" + return rel, nil + } + + // If Replace is true, we need to supercede the last release. + if i.Replace { + if err := i.replaceRelease(rel); err != nil { + return nil, err + } + } + + sortedResources, err := GetInstallSequence(chrt, sortedManifest, resources) + if err != nil { + rel.SetStatus(release.StatusFailed, fmt.Sprintf("failed to get install sequence: %s", err.Error())) + return rel, err + } + + // Store the release in history before continuing (new in Helm 3). We always know + // that this is a create operation. + if err := i.cfg.Releases.Create(rel); err != nil { + // We could try to recover gracefully here, but since nothing has been installed + // yet, this is probably safer than trying to continue when we know storage is + // not working. + return rel, err + } + rChan := make(chan resultMessage) + ctxChan := make(chan resultMessage) + doneChan := make(chan struct{}) + defer close(doneChan) + go i.performCustomInstall(rChan, rel, toBeAdopted, sortedResources) + go i.handleContext(ctx, ctxChan, doneChan, rel) + select { + case result := <-rChan: + return result.r, result.e + case result := <-ctxChan: + return result.r, result.e + } +} + func (i *Install) performInstall(c chan<- resultMessage, rel *release.Release, toBeAdopted kube.ResourceList, resources kube.ResourceList) { // pre-install hooks @@ -452,6 +933,112 @@ func (i *Install) performInstall(c chan<- resultMessage, rel *release.Release, t i.reportToRun(c, rel, nil) } + +func (i *Install) performCustomInstall(c chan<- resultMessage, rel *release.Release, toBeAdopted kube.ResourceList, installItems []InstallItem) { + + // pre-install hooks + if !i.DisableHooks { + if err := i.cfg.execHook(rel, release.HookPreInstall, i.Timeout); err != nil { + i.reportToRun(c, rel, fmt.Errorf("failed pre-install: %s", err)) + return + } + } + + for _, item := range installItems { + i.cfg.Log("-------------Installing chart: %s--------------", item.ChartName) + fmt.Fprintf(os.Stdout, "============Installing chart: %s\n===========", item.ChartName) + + if len(item.Resources) > 0 { + // 找出这个 item 中需要 adopt 的资源 + var itemToBeAdopted kube.ResourceList + var itemToCreate kube.ResourceList + + for _, r := range item.Resources { + adopted := false + for _, a := range toBeAdopted { + if a.Name == r.Name && a.Namespace == r.Namespace && + a.Mapping.GroupVersionKind == r.Mapping.GroupVersionKind { + adopted = true + itemToBeAdopted = append(itemToBeAdopted, r) + break + } + } + if !adopted { + itemToCreate = append(itemToCreate, r) + } + } + + // 创建新资源 + if len(itemToCreate) > 0 { + if _, err := i.cfg.KubeClient.Create(itemToCreate); err != nil { + i.reportToRun(c, rel, err) + return + } + } + + // 更新已存在的资源 + if len(itemToBeAdopted) > 0 { + if _, err := i.cfg.KubeClient.Update(itemToBeAdopted, itemToBeAdopted, i.Force); err != nil { + i.reportToRun(c, rel, err) + return + } + } + } + + // Helm 自带的 wait + if i.Wait && len(item.Resources) > 0 { + if i.WaitForJobs { + if err := i.cfg.KubeClient.WaitWithJobs(item.Resources, i.Timeout); err != nil { + i.reportToRun(c, rel, err) + return + } + } else { + if err := i.cfg.KubeClient.Wait(item.Resources, i.Timeout); err != nil { + i.reportToRun(c, rel, err) + return + } + } + } + + // 自定义 waitFor + if item.WaitFor != "" { + fmt.Fprintf(os.Stdout, "Executing waitFor: %s\n==========", item.WaitFor) + + cmd := exec.Command("bash", "-c", item.WaitFor) + stdoutStderr, err := cmd.CombinedOutput() + + if err != nil { + fmt.Fprintf(os.Stdout, "Wait command failed: %s, output: %s", err.Error(), string(stdoutStderr)) + rel.SetStatus(release.StatusFailed, fmt.Sprintf("wait failed for %s: %s", item.ChartName, string(stdoutStderr))) + i.reportToRun(c, rel, fmt.Errorf("wait for %s failed: %s", item.ChartName, string(stdoutStderr))) + return + } + + i.cfg.Log("Wait completed for %s", item.ChartName) + fmt.Fprintf(os.Stdout, "Wait completed for %s", item.ChartName) + } + } + + // post-install hooks - 只执行一次 + if !i.DisableHooks { + if err := i.cfg.execHook(rel, release.HookPostInstall, i.Timeout); err != nil { + i.reportToRun(c, rel, fmt.Errorf("failed post-install: %s", err)) + return + } + } + + if len(i.Description) > 0 { + rel.SetStatus(release.StatusDeployed, i.Description) + } else { + rel.SetStatus(release.StatusDeployed, "Install complete") + } + + if err := i.recordRelease(rel); err != nil { + i.cfg.Log("failed to record the release: %s", err) + } + + i.reportToRun(c, rel, nil) +} func (i *Install) handleContext(ctx context.Context, c chan<- resultMessage, done chan struct{}, rel *release.Release) { select { case <-ctx.Done(): diff --git a/pkg/action/install_sequence.go b/pkg/action/install_sequence.go new file mode 100644 index 000000000..c212b2857 --- /dev/null +++ b/pkg/action/install_sequence.go @@ -0,0 +1,152 @@ +package action + +import ( + "fmt" + "os" + "strings" + + "helm.sh/helm/v3/pkg/chart" + "helm.sh/helm/v3/pkg/kube" + "helm.sh/helm/v3/pkg/releaseutil" +) + +func GetInstallSequence(ch *chart.Chart, manifests []releaseutil.Manifest, resources kube.ResourceList) ([]InstallItem, error) { + + if len(ch.Metadata.Install) == 0 { + return []InstallItem{ + { + ChartName: ch.Name(), + WaitFor: "", + Resources: resources, + }, + }, nil + } + + // 构建 manifest content -> chart name 的映射 + // 用资源的 Kind+Name 作为 key 来匹配 + type resourceKey struct { + kind string + name string + } + manifestChartMap := make(map[resourceKey]string) + + for _, m := range manifests { + + if m.Head == nil || m.Head.Metadata == nil { + fmt.Fprintf(os.Stdout, "DEBUG: invalid manifest: %s in %s", m.Content, m.Name) + continue + } + + chartName := extractChartNameFromPath(m.Name) + // 从 manifest content 中解析出 kind 和 name + kind := m.Head.Kind + name := m.Head.Metadata.Name + if kind != "" && name != "" { + manifestChartMap[resourceKey{kind: kind, name: name}] = chartName + } + } + + // 构建 waitFor 映射 + waitForMap := make(map[string]string) + for _, inst := range ch.Metadata.Install { + waitForMap[inst.Name] = inst.WaitFor + } + + // 按 chart 分组 + groupMap := make(map[string]*InstallItem) + var order []string + + for _, inst := range ch.Metadata.Install { + groupMap[inst.Name] = &InstallItem{ + ChartName: inst.Name, + WaitFor: inst.WaitFor, + Resources: kube.ResourceList{}, + } + order = append(order, inst.Name) + } + + // 遍历 resources,通过 Kind+Name 找到对应的 chart + for _, r := range resources { + kind := r.Mapping.GroupVersionKind.Kind + name := r.Name + + key := resourceKey{kind: kind, name: name} + chartName, found := manifestChartMap[key] + if !found { + chartName = ch.Name() // fallback 到主 chart + } + + if group, exists := groupMap[chartName]; exists { + group.Resources = append(group.Resources, r) + } else { + groupMap[chartName] = &InstallItem{ + ChartName: chartName, + WaitFor: waitForMap[chartName], + Resources: kube.ResourceList{r}, + } + order = append(order, chartName) + } + } + + result := make([]InstallItem, 0, len(order)) + for _, name := range order { + if item := groupMap[name]; item != nil && len(item.Resources) > 0 { + result = append(result, *item) + } + } + + return result, nil +} + +// 从 manifest 路径中提取 chart name +func extractChartNameFromPath(filePath string) string { + parts := strings.Split(filePath, "/") + + for i, part := range parts { + if part == "charts" && i+1 < len(parts) { + return parts[i+1] + } + } + + if len(parts) >= 1 { + return parts[0] + } + + return filePath +} + +// 简单解析 manifest content 获取 kind 和 name +func parseKindAndName(content string) (kind string, name string) { + lines := strings.Split(content, "\n") + + inMetadata := false + for _, line := range lines { + trimmed := strings.TrimSpace(line) + + // 解析 kind + if strings.HasPrefix(trimmed, "kind:") { + kind = strings.TrimSpace(strings.TrimPrefix(trimmed, "kind:")) + } + + // 进入 metadata 块 + if trimmed == "metadata:" { + inMetadata = true + continue + } + + // 在 metadata 块中找 name + if inMetadata && strings.HasPrefix(trimmed, "name:") { + name = strings.TrimSpace(strings.TrimPrefix(trimmed, "name:")) + // 移除可能的引号 + name = strings.Trim(name, "\"'") + break + } + + // 离开 metadata 块 + if inMetadata && !strings.HasPrefix(line, " ") && !strings.HasPrefix(line, "\t") && trimmed != "" { + inMetadata = false + } + } + + return kind, name +} diff --git a/pkg/chart/chart.go b/pkg/chart/chart.go index a3bed63a3..8132cbf47 100644 --- a/pkg/chart/chart.go +++ b/pkg/chart/chart.go @@ -87,6 +87,27 @@ func (ch *Chart) AddDependency(charts ...*Chart) { } } +func (ch *Chart) IsSequenceInstall() bool { + if ch.Metadata.Install != nil { + return true + } + return false +} + +func (ch *Chart) ChartInstallOrder() []string { + order := make([]string, 0, len(ch.Dependencies())) + if ch.Metadata.Install != nil { + for _, item := range ch.Metadata.Install { + order = append(order, item.Name) + } + } else { + for _, item := range ch.Dependencies() { + order = append(order, item.Name()) + } + } + return order +} + // Root finds the root chart. func (ch *Chart) Root() *Chart { if ch.IsRoot() { diff --git a/pkg/chart/metadata.go b/pkg/chart/metadata.go index ae572abb7..f9c7b6a62 100644 --- a/pkg/chart/metadata.go +++ b/pkg/chart/metadata.go @@ -80,6 +80,13 @@ type Metadata struct { Dependencies []*Dependency `json:"dependencies,omitempty"` // Specifies the chart type: application or library Type string `json:"type,omitempty"` + + Install []ChartInstance `json:"install,omitempty"` +} + +type ChartInstance struct { + Name string `json:"name,omitempty"` + WaitFor string `json:"waitFor,omitempty"` } // Validate checks the metadata for known issues and sanitizes string diff --git a/pkg/kube/client.go b/pkg/kube/client.go index 7b3c803f9..572a45c45 100644 --- a/pkg/kube/client.go +++ b/pkg/kube/client.go @@ -132,6 +132,7 @@ func (c *Client) IsReachable() error { // Create creates Kubernetes resources specified in the resource list. func (c *Client) Create(resources ResourceList) (*Result, error) { c.Log("creating %d resource(s)", len(resources)) + fmt.Fprintf(os.Stdout, "creating %d resource(s)\n", len(resources)) if err := perform(resources, createResource); err != nil { return nil, err } @@ -383,6 +384,7 @@ func (c *Client) Update(original, target ResourceList, force bool) (*Result, err res := &Result{} c.Log("checking %d resources for changes", len(target)) + fmt.Fprintf(os.Stdout, "checking %d resources for changes\n", len(target)) err := target.Visit(func(info *resource.Info, err error) error { if err != nil { return err @@ -590,6 +592,7 @@ func batchPerform(infos ResourceList, fn func(*resource.Info) error, errs chan<- } func createResource(info *resource.Info) error { + fmt.Fprintf(os.Stdout, "creating resource %s\n", info.Mapping.GroupVersionKind.Kind) obj, err := resource.NewHelper(info.Client, info.Mapping).WithFieldManager(getManagedFieldsManager()).Create(info.Namespace, true, info.Object) if err != nil { return err @@ -668,6 +671,7 @@ func updateResource(c *Client, target *resource.Info, currentObj runtime.Object, return errors.Wrap(err, "failed to replace object") } c.Log("Replaced %q with kind %s for kind %s", target.Name, currentObj.GetObjectKind().GroupVersionKind().Kind, kind) + fmt.Fprintf(os.Stdout, "Replaced %q with kind %s for kind %s", target.Name, currentObj.GetObjectKind().GroupVersionKind().Kind, kind) } else { patch, patchType, err := createPatch(target, currentObj) if err != nil { @@ -685,6 +689,7 @@ func updateResource(c *Client, target *resource.Info, currentObj runtime.Object, } // send patch to server c.Log("Patch %s %q in namespace %s", kind, target.Name, target.Namespace) + fmt.Fprintf(os.Stdout, "Patch %s %q in namespace %s", kind, target.Name, target.Namespace) obj, err = helper.Patch(target.Namespace, target.Name, patchType, patch, nil) if err != nil { return errors.Wrapf(err, "cannot patch %q with kind %s", target.Name, kind) diff --git a/pkg/releaseutil/chart_sorter.go b/pkg/releaseutil/chart_sorter.go new file mode 100644 index 000000000..d8ffacf85 --- /dev/null +++ b/pkg/releaseutil/chart_sorter.go @@ -0,0 +1,94 @@ +/* +Copyright The Helm Authors. + +Licensed under the Apache License, Version 2.0 (the "License"); +you may not use this file except in compliance with the License. +You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + +Unless required by applicable law or agreed to in writing, software +distributed under the License is distributed on an "AS IS" BASIS, +WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +See the License for the specific language governing permissions and +limitations under the License. +*/ + +package releaseutil + +import ( + "helm.sh/helm/v3/pkg/release" + "sort" + "strings" +) + +type ChartSortOrder []string + +func sortManifestsByChart(manifests []Manifest, ordering ChartSortOrder) []Manifest { + sort.SliceStable(manifests, func(i, j int) bool { + return lessByChart(manifests[i], manifests[j], manifests[i].Name, manifests[j].Name, ordering) + }) + + return manifests +} + +func sortHooksByChart(hooks []*release.Hook, ordering ChartSortOrder) []*release.Hook { + h := hooks + sort.SliceStable(h, func(i, j int) bool { + return lessByChart(h[i], h[j], h[i].Name, h[j].Name, ordering) + }) + + return h +} + +func lessByChart(a interface{}, b interface{}, filePathA string, filePathB string, o ChartSortOrder) bool { + ordering := make(map[string]int, len(o)) + for v, k := range o { + ordering[k] = v + } + + nameA := extractChartNameFromPath(filePathA) + nameB := extractChartNameFromPath(filePathB) + + first, aok := ordering[nameA] + second, bok := ordering[nameB] + + if !aok && !bok { + // if both are unknown then sort alphabetically by kind, keep original order if same kind + if nameA != nameB { + return nameA < nameB + } + return first < second + } + // unknown kind is last + if !aok { + return false + } + if !bok { + return true + } + // sort different kinds, keep original order if same priority + return first < second +} + +// 安全地从路径中提取 chart name +func extractChartNameFromPath(filePath string) string { + parts := strings.Split(filePath, "/") + + // 查找 "charts" 关键字,取其后一个作为 chart name + // 格式: dtc/charts/loki/templates/xxx.yaml -> loki + for i, part := range parts { + if part == "charts" && i+1 < len(parts) { + return parts[i+1] + } + } + + // 如果没有 "charts/",说明是主 chart 的资源 + // 格式可能是: dtc/templates/xxx.yaml -> dtc (主 chart) + // 或者其他格式 + if len(parts) > 0 { + return parts[0] // 返回第一个部分作为主 chart name + } + + return filePath // fallback +} diff --git a/pkg/releaseutil/manifest_sorter.go b/pkg/releaseutil/manifest_sorter.go index 413de30e2..805fe3a00 100644 --- a/pkg/releaseutil/manifest_sorter.go +++ b/pkg/releaseutil/manifest_sorter.go @@ -111,6 +111,45 @@ func SortManifests(files map[string]string, apis chartutil.VersionSet, ordering return sortHooksByKind(result.hooks, ordering), sortManifestsByKind(result.generic, ordering), nil } +func SortManifestsByChart(files map[string]string, apis chartutil.VersionSet, ordering ChartSortOrder) ([]*release.Hook, []Manifest, error) { + result := &result{} + + var sortedFilePaths []string + for filePath := range files { + sortedFilePaths = append(sortedFilePaths, filePath) + } + sort.Strings(sortedFilePaths) + + for _, filePath := range sortedFilePaths { + content := files[filePath] + + // Skip partials. We could return these as a separate map, but there doesn't + // seem to be any need for that at this time. + if strings.HasPrefix(path.Base(filePath), "_") { + continue + } + // Skip empty files and log this. + if strings.TrimSpace(content) == "" { + continue + } + + manifestFile := &manifestFile{ + entries: SplitManifests(content), + path: filePath, + apis: apis, + } + + if err := manifestFile.sort(result); err != nil { + return result.hooks, result.generic, err + } + } + + hooks := sortHooksByKind(result.hooks, InstallOrder) + manifests := sortManifestsByKind(result.generic, InstallOrder) + + return sortHooksByChart(hooks, ordering), sortManifestsByChart(manifests, ordering), nil +} + // sort takes a manifestFile object which may contain multiple resource definition // entries and sorts each entry by hook types, and saves the resulting hooks and // generic manifests (or non-hooks) to the result struct.