diff --git a/cmd/apply.go b/cmd/apply.go index 59dcc256..8f8a3970 100644 --- a/cmd/apply.go +++ b/cmd/apply.go @@ -9,7 +9,7 @@ import ( // NewApplyCmd returns apply subcmd func NewApplyCmd(globalCfg *config.GlobalImpl) *cobra.Command { - applyOptions := &config.ApplyOptions{} + applyOptions := config.NewApplyOptions() cmd := &cobra.Command{ Use: "apply", @@ -75,6 +75,8 @@ func NewApplyCmd(globalCfg *config.GlobalImpl) *cobra.Command { f.IntVar(&applyOptions.TrackTimeout, "track-timeout", 0, `Timeout in seconds for kubedog tracking (0 to use default 300s timeout)`) f.BoolVar(&applyOptions.TrackLogs, "track-logs", false, "Enable log streaming with kubedog tracking (all pods)") f.BoolVar(&applyOptions.TrackFailedLogs, "track-failed-logs", false, "Enable log streaming with kubedog tracking, but only emit logs for pods that enter a failed state. Overridden by --track-logs when both are set") + f.DurationVar(&applyOptions.TrackLogsInterval, "track-logs-interval", applyOptions.TrackLogsInterval, + "Interval between kubedog log output updates (minimum 1s)") f.IntVar(&applyOptions.HelmStuckGrace, "helm-stuck-grace", 0, "When using --track-mode kubedog: if the cluster confirms all tracked resources have converged but the helm subprocess is still running, wait this many seconds before sending SIGINT to helm. Recovers from helm v4 hook waiter wedges. May leave the release secret in pending-install state requiring manual cleanup. 0 disables.") f.BoolVar(&applyOptions.TrackFailOnError, "track-fail-on-error", false, "Fail with non-zero exit code when kubedog tracking fails") f.StringVar(&applyOptions.Description, "description", "", `Set description for all releases. If set, overridesdescriptions in helmfile.yaml. Will be passed to "helm upgrade --description"`) diff --git a/cmd/sync.go b/cmd/sync.go index d406f9d5..079ee2a7 100644 --- a/cmd/sync.go +++ b/cmd/sync.go @@ -59,6 +59,8 @@ func NewSyncCmd(globalCfg *config.GlobalImpl) *cobra.Command { f.IntVar(&syncOptions.TrackTimeout, "track-timeout", 0, `Timeout in seconds for kubedog tracking (0 to use default 300s timeout)`) f.BoolVar(&syncOptions.TrackLogs, "track-logs", false, "Enable log streaming with kubedog tracking (all pods)") f.BoolVar(&syncOptions.TrackFailedLogs, "track-failed-logs", false, "Enable log streaming with kubedog tracking, but only emit logs for pods that enter a failed state. Overridden by --track-logs when both are set") + f.DurationVar(&syncOptions.TrackLogsInterval, "track-logs-interval", syncOptions.TrackLogsInterval, + "Interval between kubedog log output updates (minimum 1s)") f.IntVar(&syncOptions.HelmStuckGrace, "helm-stuck-grace", 0, "When using --track-mode kubedog: if the cluster confirms all tracked resources have converged but the helm subprocess is still running, wait this many seconds before sending SIGINT to helm. Recovers from helm v4 hook waiter wedges. May leave the release secret in pending-install state requiring manual cleanup. 0 disables.") f.BoolVar(&syncOptions.TrackFailOnError, "track-fail-on-error", false, "Fail with non-zero exit code when kubedog tracking fails") f.StringVar(&syncOptions.Description, "description", "", `Set description for all releases. If set, overrides descriptions in helmfile.yaml. Will be passed to "helm upgrade --description"`) diff --git a/docs/advanced-features.md b/docs/advanced-features.md index 63238b51..73266d6f 100644 --- a/docs/advanced-features.md +++ b/docs/advanced-features.md @@ -26,14 +26,19 @@ releases: Or use command-line flags: ```bash -helmfile apply --track-mode kubedog --track-timeout 300 --track-logs +helmfile apply --track-mode kubedog --track-timeout 300 --track-logs --track-logs-interval 1s ``` ### Configuration Options - **`trackMode`**: Set to `kubedog` to enable kubedog tracking, or `helm-legacy` to use Helm v4's legacy wait mode (default: `helm`) - **`trackTimeout`**: Timeout in seconds for tracking resources (default: 300) -- **`trackLogs`**: Enable real-time log streaming from tracked resources +- **`trackLogs`**: Print logs from tracked resources during deployment +- **`trackFailedLogs`**: Print collected logs only for pods that fail during deployment + +To see logs only for failed pods, set `trackFailedLogs: true` on a release or use `--track-failed-logs` with `helmfile apply` or `helmfile sync`. If both log options are enabled, logs from all tracked pods are printed. + +With kubedog tracking and either log option enabled, `helmfile apply` and `helmfile sync` print collected logs every 10 seconds by default. Use `--track-logs-interval` to change this interval (for example, `1s` above). The minimum is `1s`. Deployment progress is checked for changes every 10 seconds. ### Track Modes diff --git a/docs/configuration.md b/docs/configuration.md index 98164ae8..2a6e30b5 100644 --- a/docs/configuration.md +++ b/docs/configuration.md @@ -514,7 +514,8 @@ See [Advanced Features](advanced-features.md#resource-tracking-with-kubedog) for |-------|------|---------|-------------| | `trackMode` | string | `""` | Track mode: `helm`, `helm-legacy`, or `kubedog` | | `trackTimeout` | int | 300 | Tracking timeout in seconds | -| `trackLogs` | bool | false | Enable real-time log streaming | +| `trackLogs` | bool | false | Print logs from tracked resources during deployment (every 10 seconds by default; see [`--track-logs-interval`](advanced-features.md#resource-tracking-with-kubedog)) | +| `trackFailedLogs` | bool | false | Print collected logs only for pods that fail during deployment; overridden by `trackLogs` | | `trackKinds` | list | | Whitelist of resource kinds to track | | `skipKinds` | list | | Blacklist of resource kinds to skip | | `trackResources` | list | | Specific resources to track (objects with `kind`, `name`, `namespace`) | diff --git a/pkg/app/app.go b/pkg/app/app.go index cda12b71..2b17aeb6 100644 --- a/pkg/app/app.go +++ b/pkg/app/app.go @@ -2125,6 +2125,7 @@ Do you really want to apply? TrackTimeout: c.TrackTimeout(), TrackLogs: c.TrackLogs(), TrackFailedLogs: c.TrackFailedLogs(), + TrackLogsInterval: c.TrackLogsInterval(), HelmStuckGrace: c.HelmStuckGrace(), TrackFailOnError: c.TrackFailOnError(), Description: c.Description(), @@ -2622,6 +2623,7 @@ Do you really want to sync? TrackTimeout: c.TrackTimeout(), TrackLogs: c.TrackLogs(), TrackFailedLogs: c.TrackFailedLogs(), + TrackLogsInterval: c.TrackLogsInterval(), HelmStuckGrace: c.HelmStuckGrace(), TrackFailOnError: c.TrackFailOnError(), Description: c.Description(), diff --git a/pkg/app/app_test.go b/pkg/app/app_test.go index 6a43e301..53f2df63 100644 --- a/pkg/app/app_test.go +++ b/pkg/app/app_test.go @@ -13,6 +13,7 @@ import ( "strings" "sync" "testing" + "time" "github.com/google/go-cmp/cmp" "github.com/helmfile/vals" @@ -2556,6 +2557,7 @@ type applyConfig struct { trackMode string trackTimeout int trackLogs bool + trackLogsInterval time.Duration trackFailOnError bool // template-only options @@ -2804,6 +2806,10 @@ func (a applyConfig) TrackFailedLogs() bool { return false } +func (a applyConfig) TrackLogsInterval() time.Duration { + return a.trackLogsInterval +} + func (a applyConfig) HelmStuckGrace() int { return 0 } diff --git a/pkg/app/config.go b/pkg/app/config.go index 56238005..2a14234e 100644 --- a/pkg/app/config.go +++ b/pkg/app/config.go @@ -1,6 +1,8 @@ package app import ( + "time" + "go.uber.org/zap" "github.com/helmfile/helmfile/pkg/agent/llm" @@ -102,6 +104,7 @@ type ApplyConfigProvider interface { TrackTimeout() int TrackLogs() bool TrackFailedLogs() bool + TrackLogsInterval() time.Duration HelmStuckGrace() int TrackFailOnError() bool @@ -146,6 +149,7 @@ type SyncConfigProvider interface { TrackTimeout() int TrackLogs() bool TrackFailedLogs() bool + TrackLogsInterval() time.Duration HelmStuckGrace() int TrackFailOnError() bool diff --git a/pkg/config/apply.go b/pkg/config/apply.go index 4b4431cf..1c43e067 100644 --- a/pkg/config/apply.go +++ b/pkg/config/apply.go @@ -3,6 +3,7 @@ package config import ( "fmt" "slices" + "time" ) // ApplyOptoons is the options for the apply command @@ -95,6 +96,8 @@ type ApplyOptions struct { TrackLogs bool // TrackFailedLogs streams logs only for pods that enter a failed state. TrackFailedLogs bool + // TrackLogsInterval is the interval between kubedog log output updates. + TrackLogsInterval time.Duration // HelmStuckGrace, when > 0, enables the helm-killer safety valve. See // ReleaseSpec.HelmStuckGrace for details. Value is in seconds. HelmStuckGrace int @@ -109,7 +112,7 @@ type ApplyOptions struct { // NewApply creates a new Apply func NewApplyOptions() *ApplyOptions { - return &ApplyOptions{} + return &ApplyOptions{TrackLogsInterval: defaultTrackLogsInterval} } // ApplyImpl is impl for applyOptions @@ -346,6 +349,11 @@ func (a *ApplyImpl) TrackFailedLogs() bool { return a.ApplyOptions.TrackFailedLogs } +// TrackLogsInterval returns the kubedog log output interval. +func (a *ApplyImpl) TrackLogsInterval() time.Duration { + return a.ApplyOptions.TrackLogsInterval +} + // HelmStuckGrace returns the helm-stuck-grace value (seconds, 0 = disabled). func (a *ApplyImpl) HelmStuckGrace() int { return a.ApplyOptions.HelmStuckGrace @@ -371,5 +379,8 @@ func (a *ApplyImpl) ValidateConfig() error { if a.ApplyOptions.TrackMode != "" && !slices.Contains(validTrackModes, a.ApplyOptions.TrackMode) { return fmt.Errorf("--track-mode must be 'helm', 'helm-legacy', or 'kubedog', got: %s", a.ApplyOptions.TrackMode) } + if a.ApplyOptions.TrackLogsInterval < time.Second { + return fmt.Errorf("--track-logs-interval must be at least 1s, got: %s", a.ApplyOptions.TrackLogsInterval) + } return a.GlobalImpl.ValidateConfig() } diff --git a/pkg/config/sync.go b/pkg/config/sync.go index 83962269..318c1939 100644 --- a/pkg/config/sync.go +++ b/pkg/config/sync.go @@ -3,8 +3,14 @@ package config import ( "fmt" "slices" + "time" ) +// defaultTrackLogsInterval is the default --track-logs-interval shared by +// the sync and apply commands. It mirrors kubedog's defaultLogsInterval; +// keep the two in sync. +const defaultTrackLogsInterval = 10 * time.Second + // SyncOptions is the options for the build command type SyncOptions struct { // Set is the set flag @@ -63,6 +69,8 @@ type SyncOptions struct { TrackLogs bool // TrackFailedLogs streams logs only for pods that enter a failed state. TrackFailedLogs bool + // TrackLogsInterval is the interval between kubedog log output updates. + TrackLogsInterval time.Duration // HelmStuckGrace, when > 0, enables the helm-killer safety valve. See // ReleaseSpec.HelmStuckGrace for details. Value is in seconds. HelmStuckGrace int @@ -93,7 +101,7 @@ type SyncOptions struct { // NewSyncOptions creates a new Apply func NewSyncOptions() *SyncOptions { - return &SyncOptions{} + return &SyncOptions{TrackLogsInterval: defaultTrackLogsInterval} } // SyncImpl is impl for applyOptions @@ -254,6 +262,11 @@ func (t *SyncImpl) TrackFailedLogs() bool { return t.SyncOptions.TrackFailedLogs } +// TrackLogsInterval returns the kubedog log output interval. +func (t *SyncImpl) TrackLogsInterval() time.Duration { + return t.SyncOptions.TrackLogsInterval +} + // HelmStuckGrace returns the helm-stuck-grace value (seconds, 0 = disabled). func (t *SyncImpl) HelmStuckGrace() int { return t.SyncOptions.HelmStuckGrace @@ -349,5 +362,8 @@ func (t *SyncImpl) ValidateConfig() error { if t.SyncOptions.TrackMode != "" && !slices.Contains(validTrackModes, t.SyncOptions.TrackMode) { return fmt.Errorf("--track-mode must be 'helm', 'helm-legacy', or 'kubedog', got: %s", t.SyncOptions.TrackMode) } + if t.SyncOptions.TrackLogsInterval < time.Second { + return fmt.Errorf("--track-logs-interval must be at least 1s, got: %s", t.SyncOptions.TrackLogsInterval) + } return t.GlobalImpl.ValidateConfig() } diff --git a/pkg/config/track_logs_interval_test.go b/pkg/config/track_logs_interval_test.go new file mode 100644 index 00000000..884d457f --- /dev/null +++ b/pkg/config/track_logs_interval_test.go @@ -0,0 +1,61 @@ +package config + +import ( + "testing" + "time" + + "github.com/stretchr/testify/require" +) + +func TestTrackLogsIntervalDefault(t *testing.T) { + require.Equal(t, 10*time.Second, NewApplyOptions().TrackLogsInterval) + require.Equal(t, 10*time.Second, NewSyncOptions().TrackLogsInterval) +} + +func TestTrackLogsIntervalValidation(t *testing.T) { + commands := []struct { + name string + validate func(time.Duration) error + }{ + { + name: "apply", + validate: func(interval time.Duration) error { + options := NewApplyOptions() + options.TrackLogsInterval = interval + return NewApplyImpl(NewGlobalImpl(&GlobalOptions{}), options).ValidateConfig() + }, + }, + { + name: "sync", + validate: func(interval time.Duration) error { + options := NewSyncOptions() + options.TrackLogsInterval = interval + return NewSyncImpl(NewGlobalImpl(&GlobalOptions{}), options).ValidateConfig() + }, + }, + } + + for _, command := range commands { + t.Run(command.name, func(t *testing.T) { + for _, tt := range []struct { + name string + interval time.Duration + wantErr bool + }{ + {name: "zero", interval: 0, wantErr: true}, + {name: "negative", interval: -time.Second, wantErr: true}, + {name: "below minimum", interval: 500 * time.Millisecond, wantErr: true}, + {name: "minimum", interval: time.Second}, + } { + t.Run(tt.name, func(t *testing.T) { + err := command.validate(tt.interval) + if tt.wantErr { + require.ErrorContains(t, err, "--track-logs-interval must be at least 1s") + } else { + require.NoError(t, err) + } + }) + } + }) + } +} diff --git a/pkg/kubedog/options.go b/pkg/kubedog/options.go index 0aece172..63dbe0e8 100644 --- a/pkg/kubedog/options.go +++ b/pkg/kubedog/options.go @@ -26,6 +26,10 @@ const ( TrackModeKubedog TrackMode = "kubedog" ) +// defaultLogsInterval is how often the printer flushes captured logs when +// TrackOptions.LogsInterval is not set (zero). +const defaultLogsInterval = 10 * time.Second + type TrackOptions struct { Timeout time.Duration // Logs enables emitting logs for every pod kubedog observes. @@ -34,10 +38,12 @@ type TrackOptions struct { // them only for pods that enter a failed state (CrashLoopBackOff, Error, // ImagePullBackOff, etc.). Has no effect when Logs is true. FailedLogsOnly bool - LogsSince time.Duration - Filter *resource.FilterConfig - QPS float32 - Burst int + // LogsInterval controls how often captured logs are printed. + LogsInterval time.Duration + LogsSince time.Duration + Filter *resource.FilterConfig + QPS float32 + Burst int // Baselines holds the pre-change state of each resource keyed by // "Kind/Namespace/Name". When set, the tracker delays attaching kubedog // to a resource until its UID changes or its generation increments past @@ -51,10 +57,11 @@ type TrackOptions struct { func NewTrackOptions() *TrackOptions { return &TrackOptions{ - Timeout: 5 * time.Minute, - LogsSince: 10 * time.Minute, - QPS: 100, - Burst: 200, + Timeout: 5 * time.Minute, + LogsInterval: defaultLogsInterval, + LogsSince: 10 * time.Minute, + QPS: 100, + Burst: 200, } } @@ -68,6 +75,12 @@ func (o *TrackOptions) WithLogs(logs bool) *TrackOptions { return o } +// WithLogsInterval sets how often captured logs are printed. +func (o *TrackOptions) WithLogsInterval(interval time.Duration) *TrackOptions { + o.LogsInterval = interval + return o +} + func (o *TrackOptions) WithFilterConfig(config *resource.FilterConfig) *TrackOptions { o.Filter = config return o diff --git a/pkg/kubedog/printer.go b/pkg/kubedog/printer.go index e800beb4..63a11d0b 100644 --- a/pkg/kubedog/printer.go +++ b/pkg/kubedog/printer.go @@ -19,7 +19,6 @@ import ( const ( progressInterval = 10 * time.Second - logsInterval = 10 * time.Second heartbeatInterval = 2 * time.Minute ) @@ -191,6 +190,7 @@ type progressPrinter struct { gates *gateStatuses skipped *skippedKeys useColor bool + logsInterval time.Duration startTime time.Time // lastEmit tracks the most recent time we printed *anything* to the // logger (progress block or log stream). The heartbeat ticker uses it to @@ -206,6 +206,9 @@ type progressPrinter struct { lastLogSource string } +// newProgressPrinter builds the printer that renders progress blocks, log +// streams, and heartbeats. logsInterval sets the log-flush cadence; a zero +// value falls back to defaultLogsInterval. func newProgressPrinter( logger *zap.SugaredLogger, releaseName string, @@ -216,7 +219,11 @@ func newProgressPrinter( gates *gateStatuses, skipped *skippedKeys, useColor bool, + logsInterval time.Duration, ) *progressPrinter { + if logsInterval <= 0 { + logsInterval = defaultLogsInterval + } return &progressPrinter{ logger: logger, releaseName: releaseName, @@ -227,6 +234,7 @@ func newProgressPrinter( gates: gates, skipped: skipped, useColor: useColor, + logsInterval: logsInterval, startTime: time.Now(), lastEmit: time.Now(), lastStatus: make(map[string]string), @@ -237,7 +245,7 @@ func newProgressPrinter( func (p *progressPrinter) run(ctx context.Context, done <-chan struct{}) { progressTicker := time.NewTicker(progressInterval) defer progressTicker.Stop() - logsTicker := time.NewTicker(logsInterval) + logsTicker := time.NewTicker(p.logsInterval) defer logsTicker.Stop() heartbeatTicker := time.NewTicker(heartbeatInterval) defer heartbeatTicker.Stop() @@ -248,9 +256,6 @@ func (p *progressPrinter) run(ctx context.Context, done <-chan struct{}) { select { case <-progressTicker.C: p.flushProgress() - if !p.skipLogs { - p.flushLogs() - } case <-logsTicker.C: if !p.skipLogs { p.flushLogs() @@ -258,10 +263,6 @@ func (p *progressPrinter) run(ctx context.Context, done <-chan struct{}) { case <-heartbeatTicker.C: p.flushHeartbeat() case <-done: - p.flushProgress() - if !p.skipLogs { - p.flushLogs() - } return case <-ctx.Done(): return diff --git a/pkg/kubedog/printer_test.go b/pkg/kubedog/printer_test.go index cf59205d..368b70d2 100644 --- a/pkg/kubedog/printer_test.go +++ b/pkg/kubedog/printer_test.go @@ -211,7 +211,7 @@ func TestFlushProgress_HidesStaleReadyPodWithoutAttr(t *testing.T) { taskStore.RWTransaction(func(s *statestore.TaskStore) { s.AddReadinessTaskState(taskC) }) logger, buf, mu := newBufferedLogger(t) - p := newProgressPrinter(logger, "", taskStore, logStore, true /*skipLogs*/, false, newGateStatuses(), newSkippedKeys(), false) + p := newProgressPrinter(logger, "", taskStore, logStore, true /*skipLogs*/, false, newGateStatuses(), newSkippedKeys(), false, 0) p.flushProgress() out := capturedOutput(buf, mu) @@ -238,7 +238,7 @@ func TestFlushProgress_PreReadyPhaseAttributeOverridesReady(t *testing.T) { }) logger, buf, mu := newBufferedLogger(t) - p := newProgressPrinter(logger, "", taskStore, logStore, true, false, newGateStatuses(), newSkippedKeys(), false) + p := newProgressPrinter(logger, "", taskStore, logStore, true, false, newGateStatuses(), newSkippedKeys(), false, 0) p.flushProgress() out := capturedOutput(buf, mu) @@ -266,7 +266,7 @@ func TestFlushProgress_KeepsNamespacesWhenReleaseSpansMultiple(t *testing.T) { }) logger, buf, mu := newBufferedLogger(t) - p := newProgressPrinter(logger, "vray", taskStore, logStore, true, false, newGateStatuses(), newSkippedKeys(), false) + p := newProgressPrinter(logger, "vray", taskStore, logStore, true, false, newGateStatuses(), newSkippedKeys(), false, 0) p.flushProgress() out := capturedOutput(buf, mu) @@ -292,7 +292,7 @@ func TestFlushProgress_HidesEmptyNameChildPlaceholder(t *testing.T) { }) logger, buf, mu := newBufferedLogger(t) - p := newProgressPrinter(logger, "", taskStore, logStore, true, false, newGateStatuses(), newSkippedKeys(), false) + p := newProgressPrinter(logger, "", taskStore, logStore, true, false, newGateStatuses(), newSkippedKeys(), false, 0) p.flushProgress() out := capturedOutput(buf, mu) @@ -314,7 +314,7 @@ func TestFlushProgress_HeaderUsesReleaseName(t *testing.T) { }) logger, buf, mu := newBufferedLogger(t) - p := newProgressPrinter(logger, "vray", taskStore, logStore, true, false, newGateStatuses(), newSkippedKeys(), false) + p := newProgressPrinter(logger, "vray", taskStore, logStore, true, false, newGateStatuses(), newSkippedKeys(), false, 0) p.flushProgress() out := capturedOutput(buf, mu) assert.Contains(t, out, "Release 'vray' progress (", @@ -334,7 +334,7 @@ func TestFlushProgress_HeaderFallsBackWithoutReleaseName(t *testing.T) { }) logger, buf, mu := newBufferedLogger(t) - p := newProgressPrinter(logger, "", taskStore, logStore, true, false, newGateStatuses(), newSkippedKeys(), false) + p := newProgressPrinter(logger, "", taskStore, logStore, true, false, newGateStatuses(), newSkippedKeys(), false, 0) p.flushProgress() out := capturedOutput(buf, mu) assert.Contains(t, out, "kubedog progress (") @@ -356,7 +356,7 @@ func TestFlushProgress_NestingAndAlignment(t *testing.T) { }) logger, buf, mu := newBufferedLogger(t) - p := newProgressPrinter(logger, "", taskStore, logStore, true, false, newGateStatuses(), newSkippedKeys(), false) + p := newProgressPrinter(logger, "", taskStore, logStore, true, false, newGateStatuses(), newSkippedKeys(), false, 0) p.flushProgress() out := capturedOutput(buf, mu) @@ -392,7 +392,7 @@ func TestFlushProgress_RespectsGateAndSkip(t *testing.T) { skips.add(kdutil.ResourceID("skipped", "ns", deploymentGVK)) logger, buf, mu := newBufferedLogger(t) - p := newProgressPrinter(logger, "", taskStore, logStore, true, false, gates, skips, false) + p := newProgressPrinter(logger, "", taskStore, logStore, true, false, gates, skips, false, 0) p.flushProgress() out := capturedOutput(buf, mu) @@ -420,7 +420,7 @@ func TestFlushLogs_FailedOnlyMode_GatesUnchanged(t *testing.T) { addPodLogs(t, logStore, "init-bad", "bad line 1", "panic: boom") logger, buf, mu := newBufferedLogger(t) - p := newProgressPrinter(logger, "", taskStore, logStore, false /*skipLogs*/, true /*failedLogsOnly*/, newGateStatuses(), newSkippedKeys(), false) + p := newProgressPrinter(logger, "", taskStore, logStore, false /*skipLogs*/, true /*failedLogsOnly*/, newGateStatuses(), newSkippedKeys(), false, 0) p.flushLogs() out := capturedOutput(buf, mu) @@ -445,7 +445,7 @@ func TestFlushLogs_FailedOnlyMode_DoesNotAdvanceCursorForSuccess(t *testing.T) { addPodLogs(t, logStore, "flaky-pod", "early line A", "early line B") logger, buf, mu := newBufferedLogger(t) - p := newProgressPrinter(logger, "", taskStore, logStore, false, true, newGateStatuses(), newSkippedKeys(), false) + p := newProgressPrinter(logger, "", taskStore, logStore, false, true, newGateStatuses(), newSkippedKeys(), false, 0) p.flushLogs() assert.NotContains(t, capturedOutput(buf, mu), "early line", @@ -474,7 +474,7 @@ func TestFlushLogs_AllPodsMode_EmitsEverything(t *testing.T) { addPodLogs(t, logStore, "init-bad", "bad line 1") logger, buf, mu := newBufferedLogger(t) - p := newProgressPrinter(logger, "", taskStore, logStore, false, false /*failedLogsOnly=off*/, newGateStatuses(), newSkippedKeys(), false) + p := newProgressPrinter(logger, "", taskStore, logStore, false, false /*failedLogsOnly=off*/, newGateStatuses(), newSkippedKeys(), false, 0) p.flushLogs() out := capturedOutput(buf, mu) @@ -493,7 +493,7 @@ func TestFlushLogs_HeaderDedupedAcrossFlushes(t *testing.T) { }) logger, buf, mu := newBufferedLogger(t) - p := newProgressPrinter(logger, "", taskStore, logStore, false, false, newGateStatuses(), newSkippedKeys(), false) + p := newProgressPrinter(logger, "", taskStore, logStore, false, false, newGateStatuses(), newSkippedKeys(), false, 0) addPodLogs(t, logStore, "app-pod", "line 1", "line 2") p.flushLogs() @@ -625,7 +625,7 @@ func TestFlushHeartbeat_EmitsWhenIdleWithInFlight(t *testing.T) { }) logger, buf, mu := newBufferedLogger(t) - p := newProgressPrinter(logger, "vray", taskStore, logStore, true, false, newGateStatuses(), newSkippedKeys(), false) + p := newProgressPrinter(logger, "vray", taskStore, logStore, true, false, newGateStatuses(), newSkippedKeys(), false, 0) // Simulate the silent gap: pretend we last printed long enough ago that // the heartbeat ticker would now fire. p.lastEmit = time.Now().Add(-3 * heartbeatInterval) @@ -664,7 +664,7 @@ func TestFlushHeartbeat_LabelsGatedTasksAsWaiting(t *testing.T) { gates.set(BaselineKey("job", "ns", "queued-job"), "waiting for update (uid=abc gen=1)") logger, buf, mu := newBufferedLogger(t) - p := newProgressPrinter(logger, "vray", taskStore, logStore, true, false, gates, newSkippedKeys(), false) + p := newProgressPrinter(logger, "vray", taskStore, logStore, true, false, gates, newSkippedKeys(), false, 0) p.lastEmit = time.Now().Add(-3 * heartbeatInterval) p.flushHeartbeat() @@ -705,7 +705,7 @@ func TestFlushHeartbeat_ProgressingItemsListedBeforeWaiting(t *testing.T) { gates.set(BaselineKey("job", "ns", "bbb-queued"), "waiting for update (uid=b gen=1)") logger, buf, mu := newBufferedLogger(t) - p := newProgressPrinter(logger, "vray", taskStore, logStore, true, false, gates, newSkippedKeys(), false) + p := newProgressPrinter(logger, "vray", taskStore, logStore, true, false, gates, newSkippedKeys(), false, 0) p.lastEmit = time.Now().Add(-3 * heartbeatInterval) p.flushHeartbeat() @@ -738,7 +738,7 @@ func TestFlushHeartbeat_AllWaitingShowsCountOnly(t *testing.T) { } logger, buf, mu := newBufferedLogger(t) - p := newProgressPrinter(logger, "vray", taskStore, logStore, true, false, gates, newSkippedKeys(), false) + p := newProgressPrinter(logger, "vray", taskStore, logStore, true, false, gates, newSkippedKeys(), false, 0) p.lastEmit = time.Now().Add(-3 * heartbeatInterval) p.flushHeartbeat() @@ -768,7 +768,7 @@ func TestFlushHeartbeat_DoesNotUpdateLastEmit(t *testing.T) { }) logger, _, _ := newBufferedLogger(t) - p := newProgressPrinter(logger, "vray", taskStore, logStore, true, false, newGateStatuses(), newSkippedKeys(), false) + p := newProgressPrinter(logger, "vray", taskStore, logStore, true, false, newGateStatuses(), newSkippedKeys(), false, 0) pinned := time.Now().Add(-3 * heartbeatInterval) p.lastEmit = pinned @@ -789,7 +789,7 @@ func TestFlushHeartbeat_SuppressedWhenRecentEmit(t *testing.T) { }) logger, buf, mu := newBufferedLogger(t) - p := newProgressPrinter(logger, "vray", taskStore, logStore, true, false, newGateStatuses(), newSkippedKeys(), false) + p := newProgressPrinter(logger, "vray", taskStore, logStore, true, false, newGateStatuses(), newSkippedKeys(), false, 0) // lastEmit is current — we just printed something, so the heartbeat // should stay quiet. p.lastEmit = time.Now() @@ -816,7 +816,7 @@ func TestFlushHeartbeat_SuppressedWhenAllReady(t *testing.T) { }) logger, buf, mu := newBufferedLogger(t) - p := newProgressPrinter(logger, "vray", taskStore, logStore, true, false, newGateStatuses(), newSkippedKeys(), false) + p := newProgressPrinter(logger, "vray", taskStore, logStore, true, false, newGateStatuses(), newSkippedKeys(), false, 0) p.lastEmit = time.Now().Add(-3 * heartbeatInterval) p.flushHeartbeat() @@ -841,7 +841,7 @@ func TestFlushHeartbeat_CapsResourceListAndShowsOverflow(t *testing.T) { } logger, buf, mu := newBufferedLogger(t) - p := newProgressPrinter(logger, "vray", taskStore, logStore, true, false, newGateStatuses(), newSkippedKeys(), false) + p := newProgressPrinter(logger, "vray", taskStore, logStore, true, false, newGateStatuses(), newSkippedKeys(), false, 0) p.lastEmit = time.Now().Add(-3 * heartbeatInterval) p.flushHeartbeat() @@ -855,12 +855,12 @@ func TestFlushHeartbeat_CapsResourceListAndShowsOverflow(t *testing.T) { "progressing items beyond the cap must be summarized by overflow count, not listed") } -func TestProgressPrinter_FullRun_CancelsAndDrainsOnContextDone(t *testing.T) { - // Smoke test that the run loop terminates cleanly on context cancel. +func TestProgressPrinter_FullRun_StopsOnContextDone(t *testing.T) { + // The tracker performs the final drain after this loop has stopped. taskStore := kdutil.NewConcurrent(statestore.NewTaskStore()) logStore := kdutil.NewConcurrent(logstore.NewLogStore()) logger, _, _ := newBufferedLogger(t) - p := newProgressPrinter(logger, "", taskStore, logStore, true, false, newGateStatuses(), newSkippedKeys(), false) + p := newProgressPrinter(logger, "", taskStore, logStore, true, false, newGateStatuses(), newSkippedKeys(), false, 0) ctx, cancel := context.WithCancel(context.Background()) done := make(chan struct{}) @@ -877,6 +877,55 @@ func TestProgressPrinter_FullRun_CancelsAndDrainsOnContextDone(t *testing.T) { } } +func TestProgressPrinter_RunFlushesLogsAtConfiguredInterval(t *testing.T) { + taskStore := kdutil.NewConcurrent(statestore.NewTaskStore()) + logStore := kdutil.NewConcurrent(logstore.NewLogStore()) + ts := statestore.NewReadinessTaskState("app", "ns", deploymentGVK, statestore.ReadinessTaskStateOptions{}) + taskStore.RWTransaction(func(s *statestore.TaskStore) { + s.AddReadinessTaskState(kdutil.NewConcurrent(ts)) + }) + + logger, buf, mu := newBufferedLogger(t) + p := newProgressPrinter(logger, "", taskStore, logStore, false, false, newGateStatuses(), newSkippedKeys(), false, 20*time.Millisecond) + + ctx, cancel := context.WithCancel(context.Background()) + done := make(chan struct{}) + returned := make(chan struct{}) + go func() { + p.run(ctx, done) + close(returned) + }() + t.Cleanup(func() { + cancel() + select { + case <-returned: + case <-time.After(2 * time.Second): + t.Error("printer.run did not return after ctx cancel") + } + }) + + addPodLogs(t, logStore, "app-pod", "interval-line-one") + require.Eventually(t, func() bool { + return strings.Contains(capturedOutput(buf, mu), "interval-line-one") + }, 2*time.Second, 5*time.Millisecond, "first line should appear before tracking finishes") + + // Seeing the second line proves that another logs tick has occurred. + addPodLogs(t, logStore, "app-pod", "interval-line-two") + require.Eventually(t, func() bool { + return strings.Contains(capturedOutput(buf, mu), "interval-line-two") + }, 2*time.Second, 5*time.Millisecond) + + out := capturedOutput(buf, mu) + assert.Equal(t, 1, strings.Count(out, "interval-line-one"), "a later tick must not repeat earlier lines") + assert.Equal(t, 1, strings.Count(out, "interval-line-two")) + assert.Equal(t, 1, strings.Count(out, "kubedog progress ("), "log ticks must not print extra progress blocks") + select { + case <-returned: + t.Fatal("printer.run returned before tracking finished") + default: + } +} + // addPodLogs appends log lines to a pod's ResourceLogs entry in the log // store, creating the entry if needed. Uses an incrementing timestamp so // the printer's chronological sort is deterministic. diff --git a/pkg/kubedog/tracker.go b/pkg/kubedog/tracker.go index 7f51d66a..76585a40 100644 --- a/pkg/kubedog/tracker.go +++ b/pkg/kubedog/tracker.go @@ -429,15 +429,11 @@ func (t *Tracker) TrackResources(ctx context.Context, resources []*resource.Reso }() captureLogsFromTime := time.Now().Add(-t.trackOptions.LogsSince) - // Capture logs whenever either flag wants them. The kubedog tracker - // records logs into the logStore unconditionally; the failed-only mode - // is enforced at print time by the printer. + // Capture logs for either output mode. Failed-only filtering happens in + // the printer, after kubedog has populated the log store. wantLogsCapture := t.trackOptions.Logs || t.trackOptions.FailedLogsOnly ignoreLogs := !wantLogsCapture - // failedLogsOnly takes effect only when full streaming isn't already on. failedLogsOnly := t.trackOptions.FailedLogsOnly && !t.trackOptions.Logs - // skipLogsInPrinter mutes the printer entirely if neither mode is on. - skipLogsInPrinter := !wantLogsCapture type trackerEntry struct { target trackTarget @@ -456,7 +452,8 @@ func (t *Tracker) TrackResources(ctx context.Context, resources []*resource.Reso } gateStatuses := newGateStatuses() - printer := newProgressPrinter(t.logger, t.releaseName, taskStore, logStore, skipLogsInPrinter, failedLogsOnly, gateStatuses, t.skipped, t.trackOptions.Color) + printer := newProgressPrinter(t.logger, t.releaseName, taskStore, logStore, ignoreLogs, + failedLogsOnly, gateStatuses, t.skipped, t.trackOptions.Color, t.trackOptions.LogsInterval) printerDone := make(chan struct{}) // Spawn a parallel failure watchdog. It catches pods that genuinely @@ -553,7 +550,7 @@ func (t *Tracker) TrackResources(ctx context.Context, resources []*resource.Reso if t.releaseName != "" { releaseLabel = fmt.Sprintf("release '%s'", t.releaseName) } - t.logger.Warnf("kubedog tracker stopped for %s: %v — helmfile will continue waiting for helm to finish; no more progress/log output until then", + t.logger.Warnf("kubedog tracker stopped for %s: %v — helmfile will continue waiting for helm to finish; no more periodic progress/log output until then", releaseLabel, err) cancel() case <-trackersDone: @@ -564,6 +561,13 @@ func (t *Tracker) TrackResources(ctx context.Context, resources []*resource.Reso cancel() <-trackersDone <-printerDone + // Flush after all producers and the periodic printer have stopped. This + // guarantees the last lines are printed even if cancellation wins the + // printer's select over trackersDone. + printer.flushProgress() + if !ignoreLogs { + printer.flushLogs() + } // Drain remaining tracker errors so they're surfaced in logs. close(errCh) diff --git a/pkg/kubedog/tracker_test.go b/pkg/kubedog/tracker_test.go index abfb4b02..cee8b28f 100644 --- a/pkg/kubedog/tracker_test.go +++ b/pkg/kubedog/tracker_test.go @@ -33,6 +33,7 @@ func TestNewTrackOptions(t *testing.T) { assert.Equal(t, 5*time.Minute, opts.Timeout) assert.Equal(t, false, opts.Logs) assert.Equal(t, 10*time.Minute, opts.LogsSince) + assert.Equal(t, 10*time.Second, opts.LogsInterval) } func TestTrackOptions_WithTimeout(t *testing.T) { @@ -49,6 +50,13 @@ func TestTrackOptions_WithLogs(t *testing.T) { assert.True(t, opts.Logs) } +func TestTrackOptions_WithLogsInterval(t *testing.T) { + opts := NewTrackOptions() + opts = opts.WithLogsInterval(3 * time.Second) + + assert.Equal(t, 3*time.Second, opts.LogsInterval) +} + func TestTrackOptions_Chaining(t *testing.T) { opts := NewTrackOptions() opts = opts. @@ -375,7 +383,8 @@ func TestTrackOptions_Chaining_AllSetters(t *testing.T) { WithQPS(50). WithBurst(80). WithColor(true). - WithBaselines(baselines) + WithBaselines(baselines). + WithLogsInterval(30 * time.Second) assert.Equal(t, 2*time.Minute, opts.Timeout) assert.True(t, opts.Logs) @@ -385,6 +394,7 @@ func TestTrackOptions_Chaining_AllSetters(t *testing.T) { assert.Equal(t, 80, opts.Burst) assert.True(t, opts.Color) assert.Equal(t, baselines, opts.Baselines) + assert.Equal(t, 30*time.Second, opts.LogsInterval) } // makeObj builds an *unstructured.Unstructured with the given path/value diff --git a/pkg/state/helmx.go b/pkg/state/helmx.go index a0e0d606..9929bf96 100644 --- a/pkg/state/helmx.go +++ b/pkg/state/helmx.go @@ -330,6 +330,9 @@ func (st *HelmState) buildReleaseTracker(release *ReleaseSpec, opts *SyncOpts, u WithFailedLogsOnly(trackFailedLogs). WithFilterConfig(filterConfig). WithColor(useColor) + if opts != nil && opts.TrackLogsInterval > 0 { + trackOpts.WithLogsInterval(opts.TrackLogsInterval) + } tracker, err := kubedog.NewTracker(&kubedog.TrackerConfig{ Logger: st.logger, diff --git a/pkg/state/state.go b/pkg/state/state.go index 3d3f6a54..9ad9beb3 100644 --- a/pkg/state/state.go +++ b/pkg/state/state.go @@ -1083,6 +1083,7 @@ type SyncOpts struct { TrackTimeout int TrackLogs bool TrackFailedLogs bool + TrackLogsInterval time.Duration HelmStuckGrace int TrackFailOnError bool Description string