Files
helmfile/pkg/kubedog/tracker.go
T
yxxheroandcopilot-swe-agent[bot] 27015e8d53 fix: restore kubedog status progress output during tracking (#2602)
* fix: restore kubedog status progress output during tracking

The refactor in commit bda57b74 that replaced multitrack.Multitrack() with
individual resource trackers only read from Ready/Failed/Succeeded channels,
ignoring Status, Added, EventMsg, PodLogChunk, PodError, and AddedPod channels.
This caused kubedog status messages to no longer be displayed.

Additionally, IgnoreLogs was not passed to tracker.Options, so the trackLogs
setting was effectively ignored.

This fix restores the original multitrack-style table display using the same
kubedog utils.Table and indicators packages for:
- Formatted status tables with DEPLOYMENT/REPLICAS/AVAILABLE/UP-TO-DATE columns
- Pod sub-tables showing POD/READY/RESTARTS/STATUS with tree structure
- ANSI color coding (green=ready, yellow=in-progress, red=failed)
- Progress indicators showing value transitions (e.g. 1->3)
- Waiting messages in blue

Fixes #2601

Signed-off-by: yxxhero <aiopsclub@163.com>

* fix: address review feedback - caption coloring, termWidth, O(1) pod detection, display tests

Agent-Logs-Url: https://github.com/helmfile/helmfile/sessions/147fc763-c3f2-4a7e-9591-6f972fb62667

Co-authored-by: yxxhero <11087727+yxxhero@users.noreply.github.com>

* fix: use status.FailedReason for canary final display, fix test name typo

Agent-Logs-Url: https://github.com/helmfile/helmfile/sessions/147fc763-c3f2-4a7e-9591-6f972fb62667

Co-authored-by: yxxhero <11087727+yxxhero@users.noreply.github.com>

* fix: correct gci import grouping in display.go and display_test.go

Agent-Logs-Url: https://github.com/helmfile/helmfile/sessions/7e8f8219-5979-44fb-9729-6138c3aae08b

Co-authored-by: yxxhero <11087727+yxxhero@users.noreply.github.com>

* fix: force ANSI color output in display_test.go for CI non-TTY environments

Agent-Logs-Url: https://github.com/helmfile/helmfile/sessions/ff37ccd9-f4d1-4d42-a7d0-4903e2b9d253

Co-authored-by: yxxhero <11087727+yxxhero@users.noreply.github.com>

---------

Signed-off-by: yxxhero <aiopsclub@163.com>
Co-authored-by: copilot-swe-agent[bot] <198982749+Copilot@users.noreply.github.com>
2026-05-20 20:53:03 +08:00

621 lines
19 KiB
Go

package kubedog
import (
"context"
"fmt"
"math"
"os"
"strings"
"sync"
"time"
"github.com/werf/kubedog/pkg/display"
"github.com/werf/kubedog/pkg/informer"
"github.com/werf/kubedog/pkg/tracker"
"github.com/werf/kubedog/pkg/tracker/canary"
"github.com/werf/kubedog/pkg/tracker/daemonset"
"github.com/werf/kubedog/pkg/tracker/deployment"
"github.com/werf/kubedog/pkg/tracker/job"
"github.com/werf/kubedog/pkg/tracker/statefulset"
"github.com/werf/kubedog/pkg/trackers/dyntracker/util"
"go.uber.org/zap"
"k8s.io/apimachinery/pkg/api/meta"
"k8s.io/client-go/discovery"
"k8s.io/client-go/discovery/cached/memory"
"k8s.io/client-go/dynamic"
"k8s.io/client-go/kubernetes"
"k8s.io/client-go/rest"
"k8s.io/client-go/restmapper"
"k8s.io/client-go/tools/clientcmd"
watchtools "k8s.io/client-go/tools/watch"
"github.com/helmfile/helmfile/pkg/resource"
)
type cacheKey struct {
kubeContext string
kubeconfig string
qps float32
burst int
}
type clientCacheEntry struct {
clientSet kubernetes.Interface
dynamicClient dynamic.Interface
restConfig *rest.Config
discovery discovery.CachedDiscoveryInterface
mapper meta.RESTMapper
}
var (
kubeInitMu sync.Mutex
clientCache = make(map[cacheKey]clientCacheEntry)
)
type Tracker struct {
logger *zap.SugaredLogger
clientSet kubernetes.Interface
dynamicClient dynamic.Interface
discovery discovery.CachedDiscoveryInterface
mapper meta.RESTMapper
trackOptions *TrackOptions
filter *resource.ResourceFilter
namespace string
}
type TrackerConfig struct {
Logger *zap.SugaredLogger
Namespace string
KubeContext string
Kubeconfig string
TrackOptions *TrackOptions
KubedogQPS *float32
KubedogBurst *int
}
func NewTracker(config *TrackerConfig) (*Tracker, error) {
logger := config.Logger
if logger == nil {
logger = zap.NewNop().Sugar()
}
kubeconfig := config.Kubeconfig
if kubeconfig == "" {
kubeconfig = os.Getenv("KUBECONFIG")
}
options := config.TrackOptions
if options == nil {
options = NewTrackOptions()
}
qps := options.QPS
if config.KubedogQPS != nil {
qps = *config.KubedogQPS
}
burst := options.Burst
if config.KubedogBurst != nil {
burst = *config.KubedogBurst
}
if qps <= 0 || math.IsInf(float64(qps), 0) || math.IsNaN(float64(qps)) {
return nil, fmt.Errorf("invalid kubedog QPS %v: must be > 0 and finite", qps)
}
if burst < 1 {
return nil, fmt.Errorf("invalid kubedog burst %v: must be >= 1", burst)
}
cacheEntry, err := getOrCreateClients(config.KubeContext, kubeconfig, qps, burst)
if err != nil {
return nil, fmt.Errorf("failed to initialize kubernetes clients: %w", err)
}
var filter *resource.ResourceFilter
if options.Filter != nil {
filter = resource.NewResourceFilter(options.Filter, logger)
}
return &Tracker{
logger: logger,
clientSet: cacheEntry.clientSet,
dynamicClient: cacheEntry.dynamicClient,
discovery: cacheEntry.discovery,
mapper: cacheEntry.mapper,
trackOptions: options,
filter: filter,
namespace: config.Namespace,
}, nil
}
func getOrCreateClients(kubeContext, kubeconfig string, qps float32, burst int) (clientCacheEntry, error) {
key := cacheKey{
kubeContext: kubeContext,
kubeconfig: kubeconfig,
qps: qps,
burst: burst,
}
kubeInitMu.Lock()
if cache, ok := clientCache[key]; ok {
kubeInitMu.Unlock()
return cache, nil
}
kubeInitMu.Unlock()
loadingRules := clientcmd.NewDefaultClientConfigLoadingRules()
if kubeconfig != "" {
loadingRules.ExplicitPath = kubeconfig
}
overrides := &clientcmd.ConfigOverrides{}
if kubeContext != "" {
overrides.CurrentContext = kubeContext
}
cc := clientcmd.NewNonInteractiveDeferredLoadingClientConfig(loadingRules, overrides)
restConfig, err := cc.ClientConfig()
if err != nil {
return clientCacheEntry{}, fmt.Errorf("failed to load kubeconfig: %w", err)
}
restConfig.QPS = qps
restConfig.Burst = burst
clientSet, err := kubernetes.NewForConfig(restConfig)
if err != nil {
return clientCacheEntry{}, fmt.Errorf("failed to create kubernetes client: %w", err)
}
dynamicClient, err := dynamic.NewForConfig(restConfig)
if err != nil {
return clientCacheEntry{}, fmt.Errorf("failed to create dynamic client: %w", err)
}
discoveryClient := memory.NewMemCacheClient(clientSet.Discovery())
mapper := restmapper.NewDeferredDiscoveryRESTMapper(discoveryClient)
cache := clientCacheEntry{
clientSet: clientSet,
dynamicClient: dynamicClient,
restConfig: restConfig,
discovery: discoveryClient,
mapper: mapper,
}
kubeInitMu.Lock()
defer kubeInitMu.Unlock()
if existingCache, ok := clientCache[key]; ok {
return existingCache, nil
}
clientCache[key] = cache
return cache, nil
}
type trackTarget struct {
kind string
name string
namespace string
}
func (t *Tracker) TrackResources(ctx context.Context, resources []*resource.Resource) error {
if len(resources) == 0 {
t.logger.Info("No resources to track")
return nil
}
filtered := t.filterResources(resources)
if len(filtered) == 0 {
t.logger.Info("No resources to track after filtering")
return nil
}
t.logger.Infof("Tracking %d resources with kubedog (filtered from %d total)", len(filtered), len(resources))
targets := t.buildTargets(filtered)
if len(targets) == 0 {
t.logger.Info("No trackable resources found (only Deployment, StatefulSet, DaemonSet, Job, and Canary are supported)")
return nil
}
t.logger.Infof("Tracking breakdown: %s", t.targetSummary(targets))
watchErrCh := make(chan error, len(targets))
informerFactory := informer.NewConcurrentInformerFactory(
ctx.Done(),
watchErrCh,
t.dynamicClient,
informer.ConcurrentInformerFactoryOptions{},
)
opts := tracker.Options{
ParentContext: ctx,
Timeout: t.trackOptions.Timeout,
LogsFromTime: time.Now().Add(-t.trackOptions.LogsSince),
IgnoreLogs: !t.trackOptions.Logs,
}
var wg sync.WaitGroup
errCh := make(chan error, len(targets))
for _, target := range targets {
wg.Add(1)
go func(tgt trackTarget) {
defer wg.Done()
if err := t.trackSingleResource(tgt, informerFactory, opts); err != nil {
errCh <- fmt.Errorf("%s/%s tracking failed: %w", tgt.kind, tgt.name, err)
}
}(target)
}
done := make(chan struct{})
go func() {
wg.Wait()
close(done)
}()
select {
case err := <-errCh:
return fmt.Errorf("tracking failed: %w", err)
case <-done:
t.logger.Info("All resources tracked successfully")
return nil
case <-ctx.Done():
return fmt.Errorf("tracking canceled: %w", ctx.Err())
}
}
func (t *Tracker) trackSingleResource(target trackTarget, informerFactory *util.Concurrent[*informer.InformerFactory], opts tracker.Options) error {
parentContext := opts.ParentContext
if parentContext == nil {
parentContext = context.Background()
}
ctx, cancel := watchtools.ContextWithOptionalTimeout(parentContext, opts.Timeout)
defer cancel()
trackErrCh := make(chan error, 1)
doneCh := make(chan struct{})
switch target.kind {
case "deploy":
tr := deployment.NewTracker(target.name, target.namespace, t.clientSet, informerFactory, opts)
go t.runDeploymentTracker(ctx, tr, trackErrCh, doneCh)
return t.waitDeploymentTracker(ctx, tr, trackErrCh, doneCh)
case "sts":
tr := statefulset.NewTracker(target.name, target.namespace, t.clientSet, informerFactory, opts)
go t.runStatefulSetTracker(ctx, tr, trackErrCh, doneCh)
return t.waitStatefulSetTracker(ctx, tr, trackErrCh, doneCh)
case "ds":
tr := daemonset.NewTracker(target.name, target.namespace, t.clientSet, informerFactory, opts)
go t.runDaemonSetTracker(ctx, tr, trackErrCh, doneCh)
return t.waitDaemonSetTracker(ctx, tr, trackErrCh, doneCh)
case "job":
tr := job.NewTracker(target.name, target.namespace, t.clientSet, informerFactory, opts)
go t.runJobTracker(ctx, tr, trackErrCh, doneCh)
return t.waitJobTracker(ctx, tr, trackErrCh, doneCh)
case "canary":
tr := canary.NewTracker(target.name, target.namespace, t.clientSet, t.dynamicClient, informerFactory, opts)
go t.runCanaryTracker(ctx, tr, trackErrCh, doneCh)
return t.waitCanaryTracker(ctx, tr, trackErrCh, doneCh)
default:
return fmt.Errorf("unsupported resource kind: %s", target.kind)
}
}
func (t *Tracker) runDeploymentTracker(ctx context.Context, tr *deployment.Tracker, errCh chan<- error, doneCh chan<- struct{}) {
if err := tr.Track(ctx); err != nil {
errCh <- err
} else {
close(doneCh)
}
}
func (t *Tracker) waitDeploymentTracker(ctx context.Context, tr *deployment.Tracker, trackErrCh <-chan error, doneCh <-chan struct{}) error {
var prevStatus deployment.DeploymentStatus
out := statusOutput()
resourceName := fmt.Sprintf("deploy/%s", tr.ResourceName)
for {
select {
case status := <-tr.Added:
displayDeploymentStatusProgress(out, formatResourceCaption(resourceName, false, false), status, &prevStatus)
prevStatus = status
case <-tr.Ready:
displayDeploymentStatusProgress(out, formatResourceCaption(resourceName, true, false), prevStatus, &prevStatus)
t.logger.Infof("Deployment %s/%s is ready", tr.Namespace, tr.ResourceName)
return nil
case status := <-tr.Failed:
displayDeploymentStatusProgress(out, formatResourceCaption(resourceName, false, true), status, &prevStatus)
return fmt.Errorf("deployment %s/%s failed: %s", tr.Namespace, tr.ResourceName, status.FailedReason)
case status := <-tr.Status:
if status.StatusGeneration > prevStatus.StatusGeneration {
displayDeploymentStatusProgress(out, formatResourceCaption(resourceName, false, false), status, &prevStatus)
prevStatus = status
}
case msg := <-tr.EventMsg:
t.logger.Infof("deploy/%s: %s", tr.ResourceName, msg)
case chunk := <-tr.PodLogChunk:
t.logPodLogChunk(chunk.PodName, chunk.LogLines)
case report := <-tr.PodError:
t.logger.Warnf("deploy/%s pod %s: %s: %s", tr.ResourceName, report.ReplicaSetPodError.PodName, report.ReplicaSetPodError.ContainerName, report.ReplicaSetPodError.Message)
case <-tr.AddedReplicaSet:
case <-tr.AddedPod:
case err := <-trackErrCh:
return err
case <-doneCh:
return nil
case <-ctx.Done():
return fmt.Errorf("tracking canceled for deployment %s/%s: %w", tr.Namespace, tr.ResourceName, ctx.Err())
}
}
}
func (t *Tracker) runStatefulSetTracker(ctx context.Context, tr *statefulset.Tracker, errCh chan<- error, doneCh chan<- struct{}) {
if err := tr.Track(ctx); err != nil {
errCh <- err
} else {
close(doneCh)
}
}
func (t *Tracker) waitStatefulSetTracker(ctx context.Context, tr *statefulset.Tracker, trackErrCh <-chan error, doneCh <-chan struct{}) error {
var prevStatus statefulset.StatefulSetStatus
out := statusOutput()
resourceName := fmt.Sprintf("sts/%s", tr.ResourceName)
for {
select {
case status := <-tr.Added:
displayStatefulSetStatusProgress(out, formatResourceCaption(resourceName, false, false), status, &prevStatus)
prevStatus = status
case <-tr.Ready:
displayStatefulSetStatusProgress(out, formatResourceCaption(resourceName, true, false), prevStatus, &prevStatus)
t.logger.Infof("StatefulSet %s/%s is ready", tr.Namespace, tr.ResourceName)
return nil
case status := <-tr.Failed:
displayStatefulSetStatusProgress(out, formatResourceCaption(resourceName, false, true), status, &prevStatus)
return fmt.Errorf("statefulset %s/%s failed: %s", tr.Namespace, tr.ResourceName, status.FailedReason)
case status := <-tr.Status:
if status.StatusGeneration > prevStatus.StatusGeneration {
displayStatefulSetStatusProgress(out, formatResourceCaption(resourceName, false, false), status, &prevStatus)
prevStatus = status
}
case msg := <-tr.EventMsg:
t.logger.Infof("sts/%s: %s", tr.ResourceName, msg)
case chunk := <-tr.PodLogChunk:
t.logPodLogChunk(chunk.PodName, chunk.LogLines)
case report := <-tr.PodError:
t.logger.Warnf("sts/%s pod %s: %s: %s", tr.ResourceName, report.ReplicaSetPodError.PodName, report.ReplicaSetPodError.ContainerName, report.ReplicaSetPodError.Message)
case <-tr.AddedPod:
case err := <-trackErrCh:
return err
case <-doneCh:
return nil
case <-ctx.Done():
return fmt.Errorf("tracking canceled for statefulset %s/%s: %w", tr.Namespace, tr.ResourceName, ctx.Err())
}
}
}
func (t *Tracker) runDaemonSetTracker(ctx context.Context, tr *daemonset.Tracker, errCh chan<- error, doneCh chan<- struct{}) {
if err := tr.Track(ctx); err != nil {
errCh <- err
} else {
close(doneCh)
}
}
func (t *Tracker) waitDaemonSetTracker(ctx context.Context, tr *daemonset.Tracker, trackErrCh <-chan error, doneCh <-chan struct{}) error {
var prevStatus daemonset.DaemonSetStatus
out := statusOutput()
resourceName := fmt.Sprintf("ds/%s", tr.ResourceName)
for {
select {
case status := <-tr.Added:
displayDaemonSetStatusProgress(out, formatResourceCaption(resourceName, false, false), status, &prevStatus)
prevStatus = status
case <-tr.Ready:
displayDaemonSetStatusProgress(out, formatResourceCaption(resourceName, true, false), prevStatus, &prevStatus)
t.logger.Infof("DaemonSet %s/%s is ready", tr.Namespace, tr.ResourceName)
return nil
case status := <-tr.Failed:
displayDaemonSetStatusProgress(out, formatResourceCaption(resourceName, false, true), status, &prevStatus)
return fmt.Errorf("daemonset %s/%s failed: %s", tr.Namespace, tr.ResourceName, status.FailedReason)
case status := <-tr.Status:
if status.StatusGeneration > prevStatus.StatusGeneration {
displayDaemonSetStatusProgress(out, formatResourceCaption(resourceName, false, false), status, &prevStatus)
prevStatus = status
}
case msg := <-tr.EventMsg:
t.logger.Infof("ds/%s: %s", tr.ResourceName, msg)
case chunk := <-tr.PodLogChunk:
t.logPodLogChunk(chunk.PodName, chunk.LogLines)
case report := <-tr.PodError:
t.logger.Warnf("ds/%s pod %s: %s: %s", tr.ResourceName, report.PodError.PodName, report.PodError.ContainerName, report.PodError.Message)
case <-tr.AddedPod:
case err := <-trackErrCh:
return err
case <-doneCh:
return nil
case <-ctx.Done():
return fmt.Errorf("tracking canceled for daemonset %s/%s: %w", tr.Namespace, tr.ResourceName, ctx.Err())
}
}
}
func (t *Tracker) runJobTracker(ctx context.Context, tr *job.Tracker, errCh chan<- error, doneCh chan<- struct{}) {
if err := tr.Track(ctx); err != nil {
errCh <- err
} else {
close(doneCh)
}
}
func (t *Tracker) waitJobTracker(ctx context.Context, tr *job.Tracker, trackErrCh <-chan error, doneCh <-chan struct{}) error {
var prevStatus job.JobStatus
out := statusOutput()
resourceName := fmt.Sprintf("job/%s", tr.ResourceName)
for {
select {
case status := <-tr.Added:
displayJobStatusProgress(out, formatResourceCaption(resourceName, false, false), status, &prevStatus)
prevStatus = status
case <-tr.Succeeded:
displayJobStatusProgress(out, formatResourceCaption(resourceName, true, false), prevStatus, &prevStatus)
t.logger.Infof("Job %s/%s succeeded", tr.Namespace, tr.ResourceName)
return nil
case status := <-tr.Failed:
displayJobStatusProgress(out, formatResourceCaption(resourceName, false, true), status, &prevStatus)
return fmt.Errorf("job %s/%s failed: %s", tr.Namespace, tr.ResourceName, status.FailedReason)
case status := <-tr.Status:
if status.StatusGeneration > prevStatus.StatusGeneration {
displayJobStatusProgress(out, formatResourceCaption(resourceName, false, false), status, &prevStatus)
prevStatus = status
}
case msg := <-tr.EventMsg:
t.logger.Infof("job/%s: %s", tr.ResourceName, msg)
case chunk := <-tr.PodLogChunk:
t.logPodLogChunk(chunk.PodName, chunk.LogLines)
case report := <-tr.PodError:
t.logger.Warnf("job/%s pod %s: %s: %s", tr.ResourceName, report.PodError.PodName, report.PodError.ContainerName, report.PodError.Message)
case <-tr.AddedPod:
case err := <-trackErrCh:
return err
case <-doneCh:
return nil
case <-ctx.Done():
return fmt.Errorf("tracking canceled for job %s/%s: %w", tr.Namespace, tr.ResourceName, ctx.Err())
}
}
}
func (t *Tracker) runCanaryTracker(ctx context.Context, tr *canary.Tracker, errCh chan<- error, doneCh chan<- struct{}) {
if err := tr.Track(ctx); err != nil {
errCh <- err
} else {
close(doneCh)
}
}
func (t *Tracker) waitCanaryTracker(ctx context.Context, tr *canary.Tracker, trackErrCh <-chan error, doneCh <-chan struct{}) error {
out := statusOutput()
resourceName := fmt.Sprintf("canary/%s", tr.ResourceName)
var lastView CanaryStatusView
for {
select {
case status := <-tr.Added:
view := CanaryStatusView{
Phase: string(status.CanaryStatus.Phase),
IsFailed: status.IsFailed,
}
displayCanaryStatus(out, formatResourceCaption(resourceName, false, false), view)
lastView = view
case <-tr.Succeeded:
displayCanaryStatus(out, formatResourceCaption(resourceName, true, false), CanaryStatusView{Phase: lastView.Phase})
t.logger.Infof("Canary %s/%s succeeded", tr.Namespace, tr.ResourceName)
return nil
case status := <-tr.Failed:
displayCanaryStatus(out, formatResourceCaption(resourceName, false, true), CanaryStatusView{
Phase: status.FailedReason,
IsFailed: true,
})
return fmt.Errorf("canary %s/%s failed: %s", tr.Namespace, tr.ResourceName, status.FailedReason)
case status := <-tr.Status:
view := CanaryStatusView{
Phase: func() string {
if status.StatusIndicator != nil {
return status.StatusIndicator.Value
}
return ""
}(),
Age: status.Age,
IsFailed: status.IsFailed,
}
displayCanaryStatus(out, formatResourceCaption(resourceName, false, false), view)
lastView = view
case msg := <-tr.EventMsg:
t.logger.Infof("canary/%s: %s", tr.ResourceName, msg)
case err := <-trackErrCh:
return err
case <-doneCh:
return nil
case <-ctx.Done():
return fmt.Errorf("tracking canceled for canary %s/%s: %w", tr.Namespace, tr.ResourceName, ctx.Err())
}
}
}
func (t *Tracker) logPodLogChunk(podName string, logLines []display.LogLine) {
for _, line := range logLines {
t.logger.Infof("po/%s [%s] %s", podName, line.Timestamp, line.Message)
}
}
func (t *Tracker) buildTargets(resources []*resource.Resource) []trackTarget {
var targets []trackTarget
for _, res := range resources {
namespace := res.Namespace
if namespace == "" {
namespace = t.namespace
}
kind := ""
switch strings.ToLower(res.Kind) {
case "deployment", "deploy":
kind = "deploy"
case "statefulset", "sts":
kind = "sts"
case "daemonset", "ds":
kind = "ds"
case "job":
kind = "job"
case "canary":
kind = "canary"
default:
t.logger.Debugf("Skipping unsupported kind %s for resource %s/%s", res.Kind, namespace, res.Name)
continue
}
targets = append(targets, trackTarget{
kind: kind,
name: res.Name,
namespace: namespace,
})
}
return targets
}
func (t *Tracker) targetSummary(targets []trackTarget) string {
counts := make(map[string]int)
for _, tgt := range targets {
counts[tgt.kind]++
}
parts := make([]string, 0, len(counts))
for kind, count := range counts {
parts = append(parts, fmt.Sprintf("%ss=%d", kind, count))
}
return strings.Join(parts, ", ")
}
func (t *Tracker) filterResources(resources []*resource.Resource) []*resource.Resource {
if t.filter == nil {
return resources
}
var result []*resource.Resource
for _, res := range resources {
if t.filter.ShouldTrack(res) {
result = append(result, res)
} else {
t.logger.Debugf("Skipping resource %s/%s (kind: %s) based on configuration", res.Namespace, res.Name, res.Kind)
}
}
return result
}