feat(spec): Task 7 — sequenced upgrade action with DAG-ordered batch updates

Adds performSequencedUpgrade to upgrade.go that reuses sequencedDeployment
(from sequencing.go) to apply resources in topological DAG order when
--wait=ordered is used. Each batch calls KubeClient.Update() with matching
old resources, then waits for readiness before the next batch.

Also fixes createAndWait to call setMetadataVisitor so Helm ownership
labels/annotations are set on sequenced resources (both install and upgrade).
pull/31992/head
caretak3r 8 months ago
parent 85488eed18
commit 9b5bd5e410

@ -538,6 +538,8 @@ func (i *Install) performSequencedInstall(ctx context.Context, chrt *chart.Chart
sd := &sequencedDeployment{
cfg: i.cfg,
releaseName: rel.Name,
releaseNamespace: rel.Namespace,
disableOpenAPI: i.DisableOpenAPIValidation,
serverSideApply: i.ServerSideApply,
forceConflicts: i.ForceConflicts,

@ -80,11 +80,13 @@ func buildManifestYAML(manifests []releaseutil.Manifest) string {
return buf.String()
}
// sequencedDeployment performs ordered installation of chart resources.
// sequencedDeployment performs ordered installation or upgrade of chart resources.
// It handles the two-level DAG: first subchart ordering, then resource-group
// ordering within each chart level.
type sequencedDeployment struct {
cfg *Configuration
releaseName string
releaseNamespace string
disableOpenAPI bool
serverSideApply bool
forceConflicts bool
@ -95,6 +97,12 @@ type sequencedDeployment struct {
timeout time.Duration
readinessTimeout time.Duration
deadline time.Time // overall operation deadline
// Upgrade-specific fields. When upgradeMode is true, createAndWait delegates
// to updateAndWait which calls KubeClient.Update() instead of Create().
upgradeMode bool
currentResources kube.ResourceList // full set of old (current) resources
upgradeCSAFieldManager bool // upgrade client-side apply field manager
}
// deployChartLevel deploys all resources for a single chart level in sequenced order.
@ -202,9 +210,14 @@ func (s *sequencedDeployment) deployResourceGroupBatches(ctx context.Context, ma
return nil
}
// createAndWait creates a set of manifest resources and waits for them to be ready.
// It respects both the per-batch readiness timeout and the overall operation deadline.
// createAndWait creates (or updates, in upgrade mode) a set of manifest resources
// and waits for them to be ready. It respects both the per-batch readiness timeout
// and the overall operation deadline.
func (s *sequencedDeployment) createAndWait(ctx context.Context, manifests []releaseutil.Manifest) error {
if s.upgradeMode {
return s.updateAndWait(ctx, manifests)
}
if len(manifests) == 0 {
return nil
}
@ -218,11 +231,69 @@ func (s *sequencedDeployment) createAndWait(ctx context.Context, manifests []rel
return nil
}
if err := resources.Visit(setMetadataVisitor(s.releaseName, s.releaseNamespace, true)); err != nil {
return fmt.Errorf("setting metadata for resource batch: %w", err)
}
_, err = s.cfg.KubeClient.Create(resources, kube.ClientCreateOptionServerSideApply(s.serverSideApply, false))
if err != nil {
return fmt.Errorf("creating resource batch: %w", err)
}
return s.waitForResources(resources)
}
// updateAndWait applies an upgrade batch using KubeClient.Update() and waits for readiness.
// It matches current (old) resources by objectKey to compute the per-batch diff.
func (s *sequencedDeployment) updateAndWait(ctx context.Context, manifests []releaseutil.Manifest) error {
if len(manifests) == 0 {
return nil
}
yaml := buildManifestYAML(manifests)
target, err := s.cfg.KubeClient.Build(bytes.NewBufferString(yaml), !s.disableOpenAPI)
if err != nil {
return fmt.Errorf("building resource batch: %w", err)
}
if len(target) == 0 {
return nil
}
if err := target.Visit(setMetadataVisitor(s.releaseName, s.releaseNamespace, true)); err != nil {
return fmt.Errorf("setting metadata for resource batch: %w", err)
}
// Find the subset of current (old) resources that are represented in this batch.
// Update() will handle creates (target resources not in matchingCurrent) and
// updates (resources in both). Deletions are handled separately after all batches.
targetKeys := make(map[string]bool, len(target))
for _, r := range target {
targetKeys[objectKey(r)] = true
}
var matchingCurrent kube.ResourceList
for _, r := range s.currentResources {
if targetKeys[objectKey(r)] {
matchingCurrent = append(matchingCurrent, r)
}
}
_, err = s.cfg.KubeClient.Update(
matchingCurrent,
target,
kube.ClientUpdateOptionForceReplace(s.forceReplace),
kube.ClientUpdateOptionServerSideApply(s.serverSideApply, s.forceConflicts),
kube.ClientUpdateOptionUpgradeClientSideFieldManager(s.upgradeCSAFieldManager),
)
if err != nil {
return fmt.Errorf("updating resource batch: %w", err)
}
return s.waitForResources(target)
}
// waitForResources waits for the given resources to become ready,
// applying the per-batch and overall deadline constraints.
func (s *sequencedDeployment) waitForResources(resources kube.ResourceList) error {
// Determine effective wait timeout: min(readinessTimeout, remaining time to overall deadline)
waitTimeout := s.readinessTimeout
if !s.deadline.IsZero() {
@ -238,6 +309,7 @@ func (s *sequencedDeployment) createAndWait(ctx context.Context, manifests []rel
waitTimeout = time.Minute // safe default
}
var err error
var waiter kube.Waiter
if c, ok := s.cfg.KubeClient.(kube.InterfaceWaitOptions); ok {
waiter, err = c.GetWaiterWithOptions(s.waitStrategy, s.waitOptions...)

@ -196,7 +196,7 @@ func (u *Upgrade) RunWithContext(ctx context.Context, name string, ch chart.Char
}
u.cfg.Logger().Debug("preparing upgrade", "name", name)
currentRelease, upgradedRelease, serverSideApply, err := u.prepareUpgrade(name, chrt, vals)
currentRelease, upgradedRelease, serverSideApply, sortedManifests, err := u.prepareUpgrade(name, chrt, vals)
if err != nil {
return nil, err
}
@ -204,7 +204,12 @@ func (u *Upgrade) RunWithContext(ctx context.Context, name string, ch chart.Char
u.cfg.Releases.MaxHistory = u.MaxHistory
u.cfg.Logger().Debug("performing update", "name", name)
res, err := u.performUpgrade(ctx, currentRelease, upgradedRelease, serverSideApply)
var res *release.Release
if u.WaitStrategy == kube.OrderedWaitStrategy {
res, err = u.performSequencedUpgrade(ctx, chrt, currentRelease, upgradedRelease, sortedManifests, serverSideApply)
} else {
res, err = u.performUpgrade(ctx, currentRelease, upgradedRelease, serverSideApply)
}
if err != nil {
return res, err
}
@ -221,14 +226,16 @@ func (u *Upgrade) RunWithContext(ctx context.Context, name string, ch chart.Char
}
// prepareUpgrade builds an upgraded release for an upgrade operation.
func (u *Upgrade) prepareUpgrade(name string, chart *chartv2.Chart, vals map[string]interface{}) (*release.Release, *release.Release, bool, error) {
// When WaitStrategy is OrderedWaitStrategy, sortedManifests contains the parsed
// manifests with annotation data needed for DAG sequencing; otherwise it is nil.
func (u *Upgrade) prepareUpgrade(name string, chart *chartv2.Chart, vals map[string]interface{}) (*release.Release, *release.Release, bool, []releaseutil.Manifest, error) {
if chart == nil {
return nil, nil, false, errMissingChart
return nil, nil, false, nil, errMissingChart
}
// HideSecret must be used with dry run. Otherwise, return an error.
if !isDryRun(u.DryRunStrategy) && u.HideSecret {
return nil, nil, false, errors.New("hiding Kubernetes secrets requires a dry-run mode")
return nil, nil, false, nil, errors.New("hiding Kubernetes secrets requires a dry-run mode")
}
// finds the last non-deleted release with the given name
@ -236,19 +243,19 @@ func (u *Upgrade) prepareUpgrade(name string, chart *chartv2.Chart, vals map[str
if err != nil {
// to keep existing behavior of returning the "%q has no deployed releases" error when an existing release does not exist
if errors.Is(err, driver.ErrReleaseNotFound) {
return nil, nil, false, driver.NewErrNoDeployedReleases(name)
return nil, nil, false, nil, driver.NewErrNoDeployedReleases(name)
}
return nil, nil, false, err
return nil, nil, false, nil, err
}
lastRelease, err := releaserToV1Release(lastReleasei)
if err != nil {
return nil, nil, false, err
return nil, nil, false, nil, err
}
// Concurrent `helm upgrade`s will either fail here with `errPending` or when creating the release with "already exists". This should act as a pessimistic lock.
if lastRelease.Info.Status.IsPending() {
return nil, nil, false, errPending
return nil, nil, false, nil, errPending
}
var currentRelease *release.Release
@ -261,14 +268,14 @@ func (u *Upgrade) prepareUpgrade(name string, chart *chartv2.Chart, vals map[str
var cerr error
currentRelease, cerr = releaserToV1Release(currentReleasei)
if cerr != nil {
return nil, nil, false, err
return nil, nil, false, nil, err
}
if err != nil {
if errors.Is(err, driver.ErrNoDeployedReleases) &&
(lastRelease.Info.Status == rcommon.StatusFailed || lastRelease.Info.Status == rcommon.StatusSuperseded) {
currentRelease = lastRelease
} else {
return nil, nil, false, err
return nil, nil, false, nil, err
}
}
@ -277,11 +284,11 @@ func (u *Upgrade) prepareUpgrade(name string, chart *chartv2.Chart, vals map[str
// determine if values will be reused
vals, err = u.reuseValues(chart, currentRelease, vals)
if err != nil {
return nil, nil, false, err
return nil, nil, false, nil, err
}
if err := chartutil.ProcessDependencies(chart, vals); err != nil {
return nil, nil, false, err
return nil, nil, false, nil, err
}
// Increment revision count. This is passed to templates, and also stored on
@ -297,25 +304,33 @@ func (u *Upgrade) prepareUpgrade(name string, chart *chartv2.Chart, vals map[str
caps, err := u.cfg.getCapabilities()
if err != nil {
return nil, nil, false, err
return nil, nil, false, nil, err
}
valuesToRender, err := util.ToRenderValuesWithSchemaValidation(chart, vals, options, caps, u.SkipSchemaValidation)
if err != nil {
return nil, nil, false, err
return nil, nil, false, nil, err
}
hooks, manifestDoc, notesTxt, err := u.cfg.renderResources(chart, valuesToRender, "", "", u.SubNotes, false, false, u.PostRenderer, interactWithServer(u.DryRunStrategy), u.EnableDNS, u.HideSecret)
var sortedManifests []releaseutil.Manifest
var hooks []*release.Hook
var manifestDoc *bytes.Buffer
var notesTxt string
if u.WaitStrategy == kube.OrderedWaitStrategy {
hooks, manifestDoc, notesTxt, sortedManifests, err = u.cfg.renderResourcesWithFiles(chart, valuesToRender, "", "", u.SubNotes, false, false, u.PostRenderer, interactWithServer(u.DryRunStrategy), u.EnableDNS, u.HideSecret)
} else {
hooks, manifestDoc, notesTxt, err = u.cfg.renderResources(chart, valuesToRender, "", "", u.SubNotes, false, false, u.PostRenderer, interactWithServer(u.DryRunStrategy), u.EnableDNS, u.HideSecret)
}
if err != nil {
return nil, nil, false, err
return nil, nil, false, nil, err
}
if driver.ContainsSystemLabels(u.Labels) {
return nil, nil, false, fmt.Errorf("user supplied labels contains system reserved label name. System labels: %+v", driver.GetSystemLabels())
return nil, nil, false, nil, fmt.Errorf("user supplied labels contains system reserved label name. System labels: %+v", driver.GetSystemLabels())
}
serverSideApply, err := getUpgradeServerSideValue(u.ServerSideApply, lastRelease.ApplyMethod)
if err != nil {
return nil, nil, false, err
return nil, nil, false, nil, err
}
u.cfg.Logger().Debug("determined release apply method", slog.Bool("server_side_apply", serverSideApply), slog.String("previous_release_apply_method", lastRelease.ApplyMethod))
@ -343,7 +358,7 @@ func (u *Upgrade) prepareUpgrade(name string, chart *chartv2.Chart, vals map[str
upgradedRelease.Info.Notes = notesTxt
}
err = validateManifest(u.cfg.KubeClient, manifestDoc.Bytes(), !u.DisableOpenAPIValidation)
return currentRelease, upgradedRelease, serverSideApply, err
return currentRelease, upgradedRelease, serverSideApply, sortedManifests, err
}
func (u *Upgrade) performUpgrade(ctx context.Context, originalRelease, upgradedRelease *release.Release, serverSideApply bool) (*release.Release, error) {
@ -528,6 +543,118 @@ func (u *Upgrade) releasingUpgrade(c chan<- resultMessage, upgradedRelease *rele
u.reportToPerformUpgrade(c, upgradedRelease, nil, nil)
}
// performSequencedUpgrade deploys chart resources in DAG-ordered batches when
// --wait=ordered is used. It mirrors releasingUpgrade but uses sequencedDeployment
// to apply each batch and wait for readiness before proceeding to the next batch.
func (u *Upgrade) performSequencedUpgrade(ctx context.Context, chrt *chartv2.Chart, currentRelease, upgradedRelease *release.Release, manifests []releaseutil.Manifest, serverSideApply bool) (*release.Release, error) {
// Build the full set of current (old) resources for diff matching.
current, err := u.cfg.KubeClient.Build(bytes.NewBufferString(currentRelease.Manifest), false)
if err != nil {
if strings.Contains(err.Error(), "unable to recognize \"\": no matches for kind") {
return upgradedRelease, fmt.Errorf("current release manifest contains removed kubernetes api(s) for this "+
"kubernetes version and it is therefore unable to build the kubernetes "+
"objects for performing the diff. error from kubernetes: %w", err)
}
return upgradedRelease, fmt.Errorf("unable to build kubernetes objects from current release manifest: %w", err)
}
// Build target to compute the set of resources to delete (those removed in the new release).
target, err := u.cfg.KubeClient.Build(bytes.NewBufferString(upgradedRelease.Manifest), !u.DisableOpenAPIValidation)
if err != nil {
return upgradedRelease, fmt.Errorf("unable to build kubernetes objects from new release manifest: %w", err)
}
if isDryRun(u.DryRunStrategy) {
u.cfg.Logger().Debug("dry run for release", "name", upgradedRelease.Name)
if len(u.Description) > 0 {
upgradedRelease.Info.Description = u.Description
} else {
upgradedRelease.Info.Description = "Dry run complete"
}
return upgradedRelease, nil
}
u.cfg.Logger().Debug("creating upgraded release", "name", upgradedRelease.Name)
if err := u.cfg.Releases.Create(upgradedRelease); err != nil {
return nil, err
}
// pre-upgrade hooks
if !u.DisableHooks {
if err := u.cfg.execHook(upgradedRelease, release.HookPreUpgrade, u.WaitStrategy, u.WaitOptions, u.Timeout, serverSideApply); err != nil {
return u.failRelease(upgradedRelease, nil, fmt.Errorf("pre-upgrade hooks failed: %s", err))
}
} else {
u.cfg.Logger().Debug("upgrade hooks disabled", "name", upgradedRelease.Name)
}
upgradeCSAFieldManager := isReleaseApplyMethodClientSideApply(currentRelease.ApplyMethod) && serverSideApply
readinessTimeout := u.ReadinessTimeout
if readinessTimeout <= 0 {
readinessTimeout = time.Minute
}
sd := &sequencedDeployment{
cfg: u.cfg,
releaseName: upgradedRelease.Name,
releaseNamespace: upgradedRelease.Namespace,
disableOpenAPI: u.DisableOpenAPIValidation,
serverSideApply: serverSideApply,
forceConflicts: u.ForceConflicts,
forceReplace: u.ForceReplace,
waitStrategy: u.WaitStrategy,
waitOptions: u.WaitOptions,
waitForJobs: u.WaitForJobs,
timeout: u.Timeout,
readinessTimeout: readinessTimeout,
deadline: time.Now().Add(u.Timeout),
upgradeMode: true,
currentResources: current,
upgradeCSAFieldManager: upgradeCSAFieldManager,
}
if err := sd.deployChartLevel(ctx, chrt, manifests); err != nil {
return u.failRelease(upgradedRelease, nil, err)
}
// Delete resources that were removed in the new release (in old but not in new).
allNewKeys := make(map[string]bool, len(target))
for _, r := range target {
allNewKeys[objectKey(r)] = true
}
var toBeDeleted kube.ResourceList
for _, r := range current {
if !allNewKeys[objectKey(r)] {
toBeDeleted = append(toBeDeleted, r)
}
}
if len(toBeDeleted) > 0 {
if _, errs := u.cfg.KubeClient.Delete(toBeDeleted, metav1.DeletePropagationBackground); errs != nil {
return u.failRelease(upgradedRelease, nil, fmt.Errorf("deleting removed resources: %w", joinErrors(errs, ", ")))
}
}
// post-upgrade hooks
if !u.DisableHooks {
if err := u.cfg.execHook(upgradedRelease, release.HookPostUpgrade, u.WaitStrategy, u.WaitOptions, u.Timeout, serverSideApply); err != nil {
return u.failRelease(upgradedRelease, nil, fmt.Errorf("post-upgrade hooks failed: %s", err))
}
}
currentRelease.Info.Status = rcommon.StatusSuperseded
u.cfg.recordRelease(currentRelease)
upgradedRelease.Info.Status = rcommon.StatusDeployed
if len(u.Description) > 0 {
upgradedRelease.Info.Description = u.Description
} else {
upgradedRelease.Info.Description = "Upgrade complete"
}
u.cfg.recordRelease(upgradedRelease)
return upgradedRelease, nil
}
func (u *Upgrade) failRelease(rel *release.Release, created kube.ResourceList, err error) (*release.Release, error) {
msg := fmt.Sprintf("Upgrade %q failed: %s", rel.Name, err)
u.cfg.Logger().Warn(

@ -802,3 +802,38 @@ func TestUpgradeRelease_WaitOptionsPassedDownstream(t *testing.T) {
// Verify that WaitOptions were passed to GetWaiter
is.NotEmpty(failer.RecordedWaitOptions, "WaitOptions should be passed to GetWaiter")
}
// TestUpgradeRelease_OrderedWaitStrategy verifies that --wait=ordered upgrades
// succeed end-to-end using the fake kube client.
func TestUpgradeRelease_OrderedWaitStrategy(t *testing.T) {
req := require.New(t)
upAction := upgradeAction(t)
rel := releaseStub()
rel.Name = "seq-upgrade-test"
rel.Info.Status = common.StatusDeployed
req.NoError(upAction.cfg.Releases.Create(rel))
upAction.WaitStrategy = kube.OrderedWaitStrategy
upAction.Timeout = 5 * time.Minute
upAction.ReadinessTimeout = time.Minute
_, err := upAction.Run(rel.Name, buildChart(withSampleTemplates()), map[string]interface{}{})
req.NoError(err)
}
// TestUpgradeRelease_ReadinessTimeoutValidation checks that ReadinessTimeout > Timeout returns an error.
func TestUpgradeRelease_ReadinessTimeoutValidation(t *testing.T) {
upAction := upgradeAction(t)
rel := releaseStub()
rel.Name = "timeout-upgrade-test"
rel.Info.Status = common.StatusDeployed
require.NoError(t, upAction.cfg.Releases.Create(rel))
upAction.WaitStrategy = kube.OrderedWaitStrategy
upAction.Timeout = 5
upAction.ReadinessTimeout = 10 // exceeds Timeout
_, err := upAction.Run(rel.Name, buildChart(withSampleTemplates()), map[string]interface{}{})
require.Error(t, err)
}

Loading…
Cancel
Save