feat(dtc-install): execute post-install hooks before wait

pull/31674/head
huanghong 10 months ago
parent c372a46f32
commit fbde5f1d90

@ -107,6 +107,80 @@ func (cfg *Configuration) execHook(rl *release.Release, hook release.HookEvent,
return nil 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 // hookByWeight is a sorter for hooks
type hookByWeight []*release.Hook type hookByWeight []*release.Hook

@ -133,6 +133,8 @@ type InstallItem struct {
WaitFor string WaitFor string
Manifests []releaseutil.Manifest Manifests []releaseutil.Manifest
Resources kube.ResourceList Resources kube.ResourceList
PreInstallHooks []*release.Hook
PostInstallHooks []*release.Hook
} }
// NewInstall creates a new Install object with the given configuration. // 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 // Pre-install anything in the crd/ directory. We do this before Helm
// contacts the upstream server and builds the capabilities object. // 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 { if crds := chrt.CRDObjects(); !i.ClientOnly && !i.SkipCRDs && len(crds) > 0 {
// On dry run, bail here // On dry run, bail here
if i.DryRun { 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 { if err != nil {
rel.SetStatus(release.StatusFailed, fmt.Sprintf("failed to get install sequence: %s", err.Error())) rel.SetStatus(release.StatusFailed, fmt.Sprintf("failed to get install sequence: %s", err.Error()))
return rel, err return rel, err
@ -939,14 +942,18 @@ func (i *Install) performCustomInstall(c chan<- resultMessage, rel *release.Rele
// pre-install hooks // pre-install hooks
if !i.DisableHooks { if !i.DisableHooks {
if err := i.cfg.execHook(rel, release.HookPreInstall, i.Timeout); err != nil { 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)) i.reportToRun(c, rel, fmt.Errorf("failed pre-install: %s", err))
return return
} }
} }
for _, item := range installItems { for _, item := range installItems {
i.cfg.Log("-------------Installing chart: %s--------------", item.ChartName) i.cfg.Log("-------------Installing chart: %s--------------\n", item.ChartName)
fmt.Fprintf(os.Stdout, "============Installing chart: %s\n===========", item.ChartName) fmt.Fprintf(os.Stdout, "============ Installing chart: %s ===========\n", item.ChartName)
if len(item.Resources) > 0 { if len(item.Resources) > 0 {
// 找出这个 item 中需要 adopt 的资源 // 找出这个 item 中需要 adopt 的资源
@ -971,6 +978,11 @@ func (i *Install) performCustomInstall(c chan<- resultMessage, rel *release.Rele
// 创建新资源 // 创建新资源
if len(itemToCreate) > 0 { if len(itemToCreate) > 0 {
if _, err := i.cfg.KubeClient.Create(itemToCreate); err != nil { 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) i.reportToRun(c, rel, err)
return return
} }
@ -979,12 +991,25 @@ func (i *Install) performCustomInstall(c chan<- resultMessage, rel *release.Rele
// 更新已存在的资源 // 更新已存在的资源
if len(itemToBeAdopted) > 0 { if len(itemToBeAdopted) > 0 {
if _, err := i.cfg.KubeClient.Update(itemToBeAdopted, itemToBeAdopted, i.Force); err != nil { 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) i.reportToRun(c, rel, err)
return 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 // Helm 自带的 wait
if i.Wait && len(item.Resources) > 0 { if i.Wait && len(item.Resources) > 0 {
if i.WaitForJobs { if i.WaitForJobs {
@ -1002,31 +1027,46 @@ func (i *Install) performCustomInstall(c chan<- resultMessage, rel *release.Rele
// 自定义 waitFor // 自定义 waitFor
if item.WaitFor != "" { if item.WaitFor != "" {
fmt.Fprintf(os.Stdout, "Executing waitFor: %s\n==========", item.WaitFor) 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()))
cmd := exec.Command("bash", "-c", item.WaitFor) if err := i.recordRelease(rel); err != nil {
stdoutStderr, err := cmd.CombinedOutput() 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)))
return
} }
i.cfg.Log("Wait completed for %s", item.ChartName) i.reportToRun(c, rel, err)
fmt.Fprintf(os.Stdout, "Wait completed for %s", item.ChartName) return
} }
} }
// post-install hooks - 只执行一次 //if item.WaitFor != "" {
if !i.DisableHooks { // fmt.Fprintf(os.Stdout, "Executing waitFor: %s\n", item.WaitFor)
if err := i.cfg.execHook(rel, release.HookPostInstall, i.Timeout); err != nil { //
i.reportToRun(c, rel, fmt.Errorf("failed post-install: %s", err)) // cmd := exec.Command("bash", "-c", item.WaitFor)
return // 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 { if len(i.Description) > 0 {
rel.SetStatus(release.StatusDeployed, i.Description) rel.SetStatus(release.StatusDeployed, i.Description)
} else { } else {
@ -1039,6 +1079,66 @@ func (i *Install) performCustomInstall(c chan<- resultMessage, rel *release.Rele
i.reportToRun(c, rel, nil) 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) { func (i *Install) handleContext(ctx context.Context, c chan<- resultMessage, done chan struct{}, rel *release.Release) {
select { select {
case <-ctx.Done(): case <-ctx.Done():

@ -2,7 +2,9 @@ package action
import ( import (
"fmt" "fmt"
"helm.sh/helm/v3/pkg/release"
"os" "os"
"sort"
"strings" "strings"
"helm.sh/helm/v3/pkg/chart" "helm.sh/helm/v3/pkg/chart"
@ -10,7 +12,7 @@ import (
"helm.sh/helm/v3/pkg/releaseutil" "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 { if len(ch.Metadata.Install) == 0 {
return []InstallItem{ 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)) result := make([]InstallItem, 0, len(order))
for _, name := range order { for _, name := range order {
if item := groupMap[name]; item != nil && len(item.Resources) > 0 { if item := groupMap[name]; item != nil && len(item.Resources) > 0 {

@ -592,7 +592,7 @@ func batchPerform(infos ResourceList, fn func(*resource.Info) error, errs chan<-
} }
func createResource(info *resource.Info) error { 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) obj, err := resource.NewHelper(info.Client, info.Mapping).WithFieldManager(getManagedFieldsManager()).Create(info.Namespace, true, info.Object)
if err != nil { if err != nil {
return err return err
@ -671,7 +671,7 @@ func updateResource(c *Client, target *resource.Info, currentObj runtime.Object,
return errors.Wrap(err, "failed to replace 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) 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 { } else {
patch, patchType, err := createPatch(target, currentObj) patch, patchType, err := createPatch(target, currentObj)
if err != nil { if err != nil {
@ -689,7 +689,7 @@ func updateResource(c *Client, target *resource.Info, currentObj runtime.Object,
} }
// send patch to server // send patch to server
c.Log("Patch %s %q in namespace %s", kind, target.Name, target.Namespace) 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) obj, err = helper.Patch(target.Namespace, target.Name, patchType, patch, nil)
if err != nil { if err != nil {
return errors.Wrapf(err, "cannot patch %q with kind %s", target.Name, kind) return errors.Wrapf(err, "cannot patch %q with kind %s", target.Name, kind)

Loading…
Cancel
Save