From fbde5f1d90f4032083dbc566e8fc5f271b11362b Mon Sep 17 00:00:00 2001 From: huanghong Date: Tue, 2 Dec 2025 14:46:00 +0800 Subject: [PATCH] feat(dtc-install): execute post-install hooks before wait --- pkg/action/hooks.go | 74 ++++++++++++++++ pkg/action/install.go | 152 +++++++++++++++++++++++++++------ pkg/action/install_sequence.go | 32 ++++++- pkg/kube/client.go | 6 +- 4 files changed, 234 insertions(+), 30 deletions(-) diff --git a/pkg/action/hooks.go b/pkg/action/hooks.go index 40c1ffdb6..1f48fbcf0 100644 --- a/pkg/action/hooks.go +++ b/pkg/action/hooks.go @@ -107,6 +107,80 @@ func (cfg *Configuration) execHook(rl *release.Release, hook release.HookEvent, return nil } +func (cfg *Configuration) execChartHook(rl *release.Release, executingHook []*release.Hook, hook release.HookEvent, timeout time.Duration) error { + executingHooks := []*release.Hook{} + + executingHooks = append(executingHooks, executingHook...) + + // hooke are pre-ordered by kind, so keep order stable + sort.Stable(hookByWeight(executingHooks)) + + for _, h := range executingHooks { + // Set default delete policy to before-hook-creation + if h.DeletePolicies == nil || len(h.DeletePolicies) == 0 { + // TODO(jlegrone): Only apply before-hook-creation delete policy to run to completion + // resources. For all other resource types update in place if a + // resource with the same name already exists and is owned by the + // current release. + h.DeletePolicies = []release.HookDeletePolicy{release.HookBeforeHookCreation} + } + + if err := cfg.deleteHookByPolicy(h, release.HookBeforeHookCreation); err != nil { + return err + } + + resources, err := cfg.KubeClient.Build(bytes.NewBufferString(h.Manifest), true) + if err != nil { + return errors.Wrapf(err, "unable to build kubernetes object for %s hook %s", hook, h.Path) + } + + // Record the time at which the hook was applied to the cluster + h.LastRun = release.HookExecution{ + StartedAt: helmtime.Now(), + Phase: release.HookPhaseRunning, + } + cfg.recordRelease(rl) + + // As long as the implementation of WatchUntilReady does not panic, HookPhaseFailed or HookPhaseSucceeded + // should always be set by this function. If we fail to do that for any reason, then HookPhaseUnknown is + // the most appropriate value to surface. + h.LastRun.Phase = release.HookPhaseUnknown + + // Create hook resources + if _, err := cfg.KubeClient.Create(resources); err != nil { + h.LastRun.CompletedAt = helmtime.Now() + h.LastRun.Phase = release.HookPhaseFailed + return errors.Wrapf(err, "warning: Hook %s %s failed", hook, h.Path) + } + + // Watch hook resources until they have completed + err = cfg.KubeClient.WatchUntilReady(resources, timeout) + // Note the time of success/failure + h.LastRun.CompletedAt = helmtime.Now() + // Mark hook as succeeded or failed + if err != nil { + h.LastRun.Phase = release.HookPhaseFailed + // If a hook is failed, check the annotation of the hook to determine whether the hook should be deleted + // under failed condition. If so, then clear the corresponding resource object in the hook + if err := cfg.deleteHookByPolicy(h, release.HookFailed); err != nil { + return err + } + return err + } + h.LastRun.Phase = release.HookPhaseSucceeded + } + + // If all hooks are successful, check the annotation of each hook to determine whether the hook should be deleted + // under succeeded condition. If so, then clear the corresponding resource object in each hook + for _, h := range executingHooks { + if err := cfg.deleteHookByPolicy(h, release.HookSucceeded); err != nil { + return err + } + } + + return nil +} + // hookByWeight is a sorter for hooks type hookByWeight []*release.Hook diff --git a/pkg/action/install.go b/pkg/action/install.go index 1bca3c7c8..18d718ad0 100644 --- a/pkg/action/install.go +++ b/pkg/action/install.go @@ -129,10 +129,12 @@ type ChartPathOptions struct { } type InstallItem struct { - ChartName string - WaitFor string - Manifests []releaseutil.Manifest - Resources kube.ResourceList + ChartName string + WaitFor string + Manifests []releaseutil.Manifest + Resources kube.ResourceList + PreInstallHooks []*release.Hook + PostInstallHooks []*release.Hook } // NewInstall creates a new Install object with the given configuration. @@ -708,6 +710,7 @@ func (i *Install) RunWithCustomContext(ctx context.Context, chrt *chart.Chart, v // Pre-install anything in the crd/ directory. We do this before Helm // contacts the upstream server and builds the capabilities object. + fmt.Fprintf(os.Stdout, "Pre-install all crds in crds/...\n") if crds := chrt.CRDObjects(); !i.ClientOnly && !i.SkipCRDs && len(crds) > 0 { // On dry run, bail here if i.DryRun { @@ -840,7 +843,7 @@ func (i *Install) RunWithCustomContext(ctx context.Context, chrt *chart.Chart, v } } - sortedResources, err := GetInstallSequence(chrt, sortedManifest, resources) + sortedResources, err := GetInstallSequence(chrt, sortedManifest, resources, rel.Hooks) if err != nil { rel.SetStatus(release.StatusFailed, fmt.Sprintf("failed to get install sequence: %s", err.Error())) return rel, err @@ -939,14 +942,18 @@ func (i *Install) performCustomInstall(c chan<- resultMessage, rel *release.Rele // pre-install hooks if !i.DisableHooks { if err := i.cfg.execHook(rel, release.HookPreInstall, i.Timeout); err != nil { + rel.SetStatus(release.StatusFailed, fmt.Sprintf("failed pre-install: %s", err)) + if err := i.recordRelease(rel); err != nil { + i.cfg.Log("failed to record the release: %s", err) + } 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) + i.cfg.Log("-------------Installing chart: %s--------------\n", item.ChartName) + fmt.Fprintf(os.Stdout, "============ Installing chart: %s ===========\n", item.ChartName) if len(item.Resources) > 0 { // 找出这个 item 中需要 adopt 的资源 @@ -971,6 +978,11 @@ func (i *Install) performCustomInstall(c chan<- resultMessage, rel *release.Rele // 创建新资源 if len(itemToCreate) > 0 { if _, err := i.cfg.KubeClient.Create(itemToCreate); err != nil { + rel.SetStatus(release.StatusFailed, fmt.Sprintf("create failed: %s", err.Error())) + if err := i.recordRelease(rel); err != nil { + i.cfg.Log("failed to record the release: %s", err) + } + i.reportToRun(c, rel, err) return } @@ -979,12 +991,25 @@ func (i *Install) performCustomInstall(c chan<- resultMessage, rel *release.Rele // 更新已存在的资源 if len(itemToBeAdopted) > 0 { if _, err := i.cfg.KubeClient.Update(itemToBeAdopted, itemToBeAdopted, i.Force); err != nil { + rel.SetStatus(release.StatusFailed, fmt.Sprintf("update failed: %s", err.Error())) + if err := i.recordRelease(rel); err != nil { + i.cfg.Log("failed to record the release: %s", err) + } + i.reportToRun(c, rel, err) return } } } + // Post-install hooks + if !i.DisableHooks { + if err := i.cfg.execChartHook(rel, item.PostInstallHooks, release.HookPostInstall, i.Timeout); err != nil { + i.reportToRun(c, rel, fmt.Errorf("failed post-install: %s", err)) + return + } + } + // Helm 自带的 wait if i.Wait && len(item.Resources) > 0 { if i.WaitForJobs { @@ -1002,30 +1027,45 @@ func (i *Install) performCustomInstall(c chan<- resultMessage, rel *release.Rele // 自定义 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 := i.executeWaitFor(item.ChartName, item.WaitFor); err != nil { + rel.SetStatus(release.StatusFailed, fmt.Sprintf("wait failed for %s: %s", item.ChartName, err.Error())) + if err := i.recordRelease(rel); err != nil { + i.cfg.Log("failed to record the release: %s", err) + } - 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))) + i.reportToRun(c, rel, err) 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 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)) + // if strings.Contains(string(stdoutStderr), "no matching resources found") { + // goto retryWaitFor + // } + // 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\n", 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) @@ -1039,6 +1079,66 @@ func (i *Install) performCustomInstall(c chan<- resultMessage, rel *release.Rele i.reportToRun(c, rel, nil) } + +// executeWaitFor 执行自定义等待命令,带重试机制 +func (i *Install) executeWaitFor(chartName, waitFor string) error { + const ( + maxRetries = 100 + retryInterval = 5 * time.Second + ) + + // 可重试的错误关键字 + retryableErrors := []string{ + "no matching resources found", + "etcdserver: leader changed", + "etcdserver: request timed out", + "connection refused", + "connection reset by peer", + "i/o timeout", + "net/http: TLS handshake timeout", + "the object has been modified", + "unable to connect to the server", + "context deadline exceeded", + "EOF", + } + + fmt.Fprintf(os.Stdout, "Executing waitFor: %s\n", waitFor) + + for attempt := 1; attempt <= maxRetries; attempt++ { + cmd := exec.Command("bash", "-c", waitFor) + stdoutStderr, err := cmd.CombinedOutput() + + if err == nil { + i.cfg.Log("Wait completed for %s", chartName) + fmt.Fprintf(os.Stdout, "Wait completed for %s\n", chartName) + return nil + } + + output := string(stdoutStderr) + + // 检查是否是可重试的错误 + shouldRetry := false + for _, retryableErr := range retryableErrors { + if strings.Contains(output, retryableErr) { + shouldRetry = true + break + } + } + + if shouldRetry { + fmt.Fprintf(os.Stdout, "Retryable error, retrying (%d/%d): %s\n", attempt, maxRetries, strings.TrimSpace(output)) + time.Sleep(retryInterval) + continue + } + + // 不可重试的错误,直接返回 + fmt.Fprintf(os.Stdout, "Wait command failed: %s, output: %s\n", err.Error(), output) + return fmt.Errorf("wait for %s failed: %s", chartName, output) + } + + return fmt.Errorf("wait for %s timed out after %d retries", chartName, maxRetries) +} + 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 index c212b2857..ab36690f2 100644 --- a/pkg/action/install_sequence.go +++ b/pkg/action/install_sequence.go @@ -2,7 +2,9 @@ package action import ( "fmt" + "helm.sh/helm/v3/pkg/release" "os" + "sort" "strings" "helm.sh/helm/v3/pkg/chart" @@ -10,7 +12,7 @@ import ( "helm.sh/helm/v3/pkg/releaseutil" ) -func GetInstallSequence(ch *chart.Chart, manifests []releaseutil.Manifest, resources kube.ResourceList) ([]InstallItem, error) { +func GetInstallSequence(ch *chart.Chart, manifests []releaseutil.Manifest, resources kube.ResourceList, hooks []*release.Hook) ([]InstallItem, error) { if len(ch.Metadata.Install) == 0 { return []InstallItem{ @@ -88,6 +90,34 @@ func GetInstallSequence(ch *chart.Chart, manifests []releaseutil.Manifest, resou } } + for _, h := range hooks { + chartName := extractChartNameFromPath(h.Path) + + item, exists := groupMap[chartName] + if !exists { + continue + } + + for _, event := range h.Events { + switch event { + case release.HookPreInstall: + item.PreInstallHooks = append(item.PreInstallHooks, h) + case release.HookPostInstall: + item.PostInstallHooks = append(item.PostInstallHooks, h) + } + } + } + + // 对每个 item 的 hooks 按 weight 排序 + for _, item := range groupMap { + sort.Slice(item.PreInstallHooks, func(i, j int) bool { + return item.PreInstallHooks[i].Weight < item.PreInstallHooks[j].Weight + }) + sort.Slice(item.PostInstallHooks, func(i, j int) bool { + return item.PostInstallHooks[i].Weight < item.PostInstallHooks[j].Weight + }) + } + result := make([]InstallItem, 0, len(order)) for _, name := range order { if item := groupMap[name]; item != nil && len(item.Resources) > 0 { diff --git a/pkg/kube/client.go b/pkg/kube/client.go index 8fe57cca6..09352b5d5 100644 --- a/pkg/kube/client.go +++ b/pkg/kube/client.go @@ -592,7 +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 '%s'", info.Mapping.GroupVersionKind.Kind, info.Name) + fmt.Fprintf(os.Stdout, "creating resource %s '%s'\n", info.Mapping.GroupVersionKind.Kind, info.Name) obj, err := resource.NewHelper(info.Client, info.Mapping).WithFieldManager(getManagedFieldsManager()).Create(info.Namespace, true, info.Object) if err != nil { return err @@ -671,7 +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) + fmt.Fprintf(os.Stdout, "Replaced %q with kind %s for kind %s\n", target.Name, currentObj.GetObjectKind().GroupVersionKind().Kind, kind) } else { patch, patchType, err := createPatch(target, currentObj) if err != nil { @@ -689,7 +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) + fmt.Fprintf(os.Stdout, "Patch %s %q in namespace %s\n", 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)