diff --git a/KUBEDOG_CONFIG.md b/KUBEDOG_CONFIG.md index 1f095f8f..d433ec48 100644 --- a/KUBEDOG_CONFIG.md +++ b/KUBEDOG_CONFIG.md @@ -172,6 +172,23 @@ releases: trackLogs: true # Show pod logs during tracking ``` +### trackLogsUntilReady + +Limit pod log output to startup while continuing readiness and failure tracking: + +```yaml +releases: + - name: my-app + trackMode: kubedog + trackTimeout: 600 + trackLogs: true + trackLogsUntilReady: true +``` + +Each pod stops printing routine logs once it first becomes ready, even while other replicas are still starting. Logs captured before readiness are flushed. The same pod does not resume routine log output if it later becomes unready or fails. Job logging is unchanged. The default is `false`, which keeps full streaming when `trackLogs` is enabled. This option does not change `trackFailedLogs` when `trackLogs` is disabled. + +The equivalent CLI flags are `--track-logs --track-logs-until-ready` with `--track-mode kubedog`. + ### trackKinds / skipKinds Control which resource types to track: diff --git a/cmd/apply.go b/cmd/apply.go index 8f8a3970..7e4b8aa5 100644 --- a/cmd/apply.go +++ b/cmd/apply.go @@ -74,6 +74,7 @@ func NewApplyCmd(globalCfg *config.GlobalImpl) *cobra.Command { f.StringVar(&applyOptions.TrackMode, "track-mode", "", "Track mode for releases: 'helm' (default), 'helm-legacy' (Helm v4 only), or 'kubedog'") 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.TrackLogsUntilReady, "track-logs-until-ready", false, "Stop pod log streaming when each pod becomes ready (requires --track-logs)") 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)") diff --git a/cmd/root_test.go b/cmd/root_test.go index ff4af7b4..91afe532 100644 --- a/cmd/root_test.go +++ b/cmd/root_test.go @@ -82,3 +82,22 @@ func TestRootCmdRegistersOtelTracingFlag(t *testing.T) { assert.Equal(t, "false", flag.DefValue) assert.Contains(t, flag.Usage, "HELMFILE_OTEL_TRACING") } + +func TestTrackLogsUntilReadyFlag(t *testing.T) { + for _, command := range []string{"sync", "apply"} { + t.Run(command, func(t *testing.T) { + rootCmd, err := NewRootCmd(&config.GlobalOptions{}) + require.NoError(t, err) + cmd, _, err := rootCmd.Find([]string{command}) + require.NoError(t, err) + + flag := cmd.Flags().Lookup("track-logs-until-ready") + require.NotNil(t, flag) + assert.Equal(t, "false", flag.DefValue) + require.NoError(t, cmd.ParseFlags([]string{"--track-logs", "--track-logs-until-ready"})) + enabled, err := cmd.Flags().GetBool("track-logs-until-ready") + require.NoError(t, err) + assert.True(t, enabled) + }) + } +} diff --git a/cmd/sync.go b/cmd/sync.go index 079ee2a7..f25a8927 100644 --- a/cmd/sync.go +++ b/cmd/sync.go @@ -58,6 +58,7 @@ func NewSyncCmd(globalCfg *config.GlobalImpl) *cobra.Command { f.StringVar(&syncOptions.TrackMode, "track-mode", "", "Track mode for releases: 'helm' (default), 'helm-legacy' (Helm v4 only), or 'kubedog'") 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.TrackLogsUntilReady, "track-logs-until-ready", false, "Stop pod log streaming when each pod becomes ready (requires --track-logs)") 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)") diff --git a/docs/advanced-features.md b/docs/advanced-features.md index 73266d6f..d3562662 100644 --- a/docs/advanced-features.md +++ b/docs/advanced-features.md @@ -34,12 +34,39 @@ helmfile apply --track-mode kubedog --track-timeout 300 --track-logs --track-log - **`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`**: Print logs from tracked resources during deployment +- **`trackLogsUntilReady`**: With `trackLogs: true`, stop routine pod logs once each pod first becomes ready (default: false) - **`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. +### Startup-Only Pod Logs + +Use `trackLogsUntilReady: true` to show startup logs without continuing to print application traffic from ready pods while other replicas are still starting: + +```yaml +releases: + - name: myapp + chart: ./charts/myapp + trackMode: kubedog + trackTimeout: 600 + trackLogs: true + trackLogsUntilReady: true +``` + +Or use command-line flags: + +```bash +helmfile apply --track-mode kubedog --track-timeout 600 --track-logs --track-logs-until-ready +``` + +Logs captured before readiness are flushed. Log output remains limited to startup if the same pod later becomes unready or fails; readiness and failure tracking continue for every pod. Job logs retain their existing behavior. + +Kubernetes readiness timestamps have second precision, so logs from the readiness second are retained to avoid losing startup output. + +The option defaults to `false`, so `trackLogs: true` continues printing logs throughout tracking. It only applies when `trackLogs` is enabled; `trackFailedLogs: true` alone keeps its existing failed-only behavior. When both `trackLogs` and `trackFailedLogs` are enabled, `trackLogs` takes precedence. + ### Track Modes Helmfile supports three track modes: diff --git a/docs/configuration.md b/docs/configuration.md index 2a6e30b5..7f81066f 100644 --- a/docs/configuration.md +++ b/docs/configuration.md @@ -515,6 +515,7 @@ 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 | Print logs from tracked resources during deployment (every 10 seconds by default; see [`--track-logs-interval`](advanced-features.md#resource-tracking-with-kubedog)) | +| `trackLogsUntilReady` | bool | false | With `trackLogs: true`, stop routine pod logs once each pod first becomes ready; readiness and failure tracking continue. Job logs are unaffected | | `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 | diff --git a/pkg/app/app.go b/pkg/app/app.go index 2b17aeb6..ac339107 100644 --- a/pkg/app/app.go +++ b/pkg/app/app.go @@ -2124,6 +2124,7 @@ Do you really want to apply? TrackMode: c.TrackMode(), TrackTimeout: c.TrackTimeout(), TrackLogs: c.TrackLogs(), + TrackLogsUntilReady: c.TrackLogsUntilReady(), TrackFailedLogs: c.TrackFailedLogs(), TrackLogsInterval: c.TrackLogsInterval(), HelmStuckGrace: c.HelmStuckGrace(), @@ -2622,6 +2623,7 @@ Do you really want to sync? TrackMode: c.TrackMode(), TrackTimeout: c.TrackTimeout(), TrackLogs: c.TrackLogs(), + TrackLogsUntilReady: c.TrackLogsUntilReady(), TrackFailedLogs: c.TrackFailedLogs(), TrackLogsInterval: c.TrackLogsInterval(), HelmStuckGrace: c.HelmStuckGrace(), diff --git a/pkg/app/app_test.go b/pkg/app/app_test.go index 53f2df63..8ec509d9 100644 --- a/pkg/app/app_test.go +++ b/pkg/app/app_test.go @@ -2557,6 +2557,7 @@ type applyConfig struct { trackMode string trackTimeout int trackLogs bool + trackLogsUntilReady bool trackLogsInterval time.Duration trackFailOnError bool @@ -2802,6 +2803,10 @@ func (a applyConfig) TrackLogs() bool { return a.trackLogs } +func (a applyConfig) TrackLogsUntilReady() bool { + return a.trackLogsUntilReady +} + func (a applyConfig) TrackFailedLogs() bool { return false } diff --git a/pkg/app/config.go b/pkg/app/config.go index 2a14234e..1b6e3400 100644 --- a/pkg/app/config.go +++ b/pkg/app/config.go @@ -103,6 +103,7 @@ type ApplyConfigProvider interface { TrackMode() string TrackTimeout() int TrackLogs() bool + TrackLogsUntilReady() bool TrackFailedLogs() bool TrackLogsInterval() time.Duration HelmStuckGrace() int @@ -148,6 +149,7 @@ type SyncConfigProvider interface { TrackMode() string TrackTimeout() int TrackLogs() bool + TrackLogsUntilReady() bool TrackFailedLogs() bool TrackLogsInterval() time.Duration HelmStuckGrace() int diff --git a/pkg/config/apply.go b/pkg/config/apply.go index 1c43e067..7477354c 100644 --- a/pkg/config/apply.go +++ b/pkg/config/apply.go @@ -94,6 +94,8 @@ type ApplyOptions struct { TrackTimeout int // TrackLogs enables log streaming with kubedog TrackLogs bool + // TrackLogsUntilReady stops pod log streaming when each pod becomes ready. Requires TrackLogs. + TrackLogsUntilReady bool // TrackFailedLogs streams logs only for pods that enter a failed state. TrackFailedLogs bool // TrackLogsInterval is the interval between kubedog log output updates. @@ -344,6 +346,11 @@ func (a *ApplyImpl) TrackLogs() bool { return a.ApplyOptions.TrackLogs } +// TrackLogsUntilReady returns the track-logs-until-ready flag. +func (a *ApplyImpl) TrackLogsUntilReady() bool { + return a.ApplyOptions.TrackLogsUntilReady +} + // TrackFailedLogs returns the track-failed-logs flag. func (a *ApplyImpl) TrackFailedLogs() bool { return a.ApplyOptions.TrackFailedLogs diff --git a/pkg/config/sync.go b/pkg/config/sync.go index 318c1939..e32c2c90 100644 --- a/pkg/config/sync.go +++ b/pkg/config/sync.go @@ -67,6 +67,8 @@ type SyncOptions struct { TrackTimeout int // TrackLogs enables log streaming with kubedog TrackLogs bool + // TrackLogsUntilReady stops pod log streaming when each pod becomes ready. Requires TrackLogs. + TrackLogsUntilReady bool // TrackFailedLogs streams logs only for pods that enter a failed state. TrackFailedLogs bool // TrackLogsInterval is the interval between kubedog log output updates. @@ -257,6 +259,11 @@ func (t *SyncImpl) TrackLogs() bool { return t.SyncOptions.TrackLogs } +// TrackLogsUntilReady returns the track-logs-until-ready flag. +func (t *SyncImpl) TrackLogsUntilReady() bool { + return t.SyncOptions.TrackLogsUntilReady +} + // TrackFailedLogs returns the track-failed-logs flag. func (t *SyncImpl) TrackFailedLogs() bool { return t.SyncOptions.TrackFailedLogs diff --git a/pkg/kubedog/logs.go b/pkg/kubedog/logs.go new file mode 100644 index 00000000..57a4752b --- /dev/null +++ b/pkg/kubedog/logs.go @@ -0,0 +1,157 @@ +package kubedog + +import ( + "fmt" + "sync" + "time" + + "github.com/werf/kubedog/pkg/informer" + kdutil "github.com/werf/kubedog/pkg/trackers/dyntracker/util" + "k8s.io/apimachinery/pkg/apis/meta/v1/unstructured" + "k8s.io/apimachinery/pkg/runtime/schema" + "k8s.io/apimachinery/pkg/types" + "k8s.io/client-go/tools/cache" +) + +type podLogCutoff struct { + uid types.UID + time time.Time +} + +// podLogFilter remembers the first Ready transition independently of kubedog's +// workload readiness. It only filters output; pod tracking and log capture continue. +type podLogFilter struct { + mu sync.RWMutex + cutoffs map[string]podLogCutoff +} + +func newPodLogFilter() *podLogFilter { + return &podLogFilter{cutoffs: make(map[string]podLogCutoff)} +} + +// newStartupLogFilter returns a readiness-tracking log filter when log +// streaming is limited to startup output, or nil otherwise. +func newStartupLogFilter(options *TrackOptions, factory *kdutil.Concurrent[*informer.InformerFactory], targets []trackTarget) (*podLogFilter, error) { + if options == nil || !options.Logs || !options.LogsUntilReady { + return nil, nil + } + filter := newPodLogFilter() + if err := filter.watch(factory, targets); err != nil { + return nil, err + } + return filter, nil +} + +// cutoffTrackedKind reports whether pods of the target kind are subject to +// readiness-based log cutoff. Jobs are excluded: they can be Ready while +// still executing. +func cutoffTrackedKind(kind string) bool { + switch kind { + case "deploy", "sts", "ds": + return true + default: + return false + } +} + +// watchNamespacePods attaches the readiness handler to the shared pod +// informer for namespace and starts it. +func (f *podLogFilter) watchNamespacePods(factory *kdutil.Concurrent[*informer.InformerFactory], gvr schema.GroupVersionResource, namespace string) error { + var podInformer *kdutil.Concurrent[*informer.Informer] + err := factory.RWTransactionErr(func(factory *informer.InformerFactory) error { + informer, err := factory.ForNamespace(gvr, namespace) + if err != nil { + return fmt.Errorf("create pod log informer: %w", err) + } + podInformer = informer + return nil + }) + if err != nil { + return err + } + + return podInformer.RWTransactionErr(func(inf *informer.Informer) error { + _, err := inf.AddEventHandler(cache.ResourceEventHandlerFuncs{ + AddFunc: f.update, + UpdateFunc: func(_, obj any) { + f.update(obj) + }, + }) + if err != nil { + return fmt.Errorf("watch pod readiness for logs: %w", err) + } + inf.Run() + return nil + }) +} + +func (f *podLogFilter) watch(factory *kdutil.Concurrent[*informer.InformerFactory], targets []trackTarget) error { + podGVR := schema.GroupVersionResource{Version: "v1", Resource: "pods"} + namespaces := make(map[string]struct{}) + for _, target := range targets { + if !cutoffTrackedKind(target.kind) { + continue + } + if _, watched := namespaces[target.namespace]; watched { + continue + } + if err := f.watchNamespacePods(factory, podGVR, target.namespace); err != nil { + return err + } + namespaces[target.namespace] = struct{}{} + } + return nil +} + +func (f *podLogFilter) update(obj any) { + pod, ok := obj.(*unstructured.Unstructured) + if !ok { + return + } + id := kdutil.ResourceID(pod.GetName(), pod.GetNamespace(), watchdogPodGVK) + f.mu.Lock() + defer f.mu.Unlock() + // Keep the first Ready transition recorded for this pod instance. A + // replacement pod reusing the name starts fresh; delete on a missing key + // is a no-op. + if cutoff, found := f.cutoffs[id]; found && cutoff.uid == pod.GetUID() { + return + } + delete(f.cutoffs, id) + // Job pods can be Ready while still executing; keep their complete output. + for _, owner := range pod.GetOwnerReferences() { + if owner.Kind == "Job" { + return + } + } + conditions, _, err := unstructured.NestedSlice(pod.Object, "status", "conditions") + if err != nil { + return + } + for _, raw := range conditions { + condition, ok := raw.(map[string]any) + if !ok || condition["type"] != "Ready" || condition["status"] != "True" { + continue + } + timestamp, _ := condition["lastTransitionTime"].(string) + readyTime, err := time.Parse(time.RFC3339Nano, timestamp) + if err != nil { + readyTime = time.Now() + } else if readyTime.Nanosecond() == 0 { + // Kubernetes rounds Ready timestamps to seconds. Preserve startup + // lines within that second too, even when delivered after readiness. + readyTime = readyTime.Add(time.Second - time.Nanosecond) + } + f.cutoffs[id] = podLogCutoff{uid: pod.GetUID(), time: readyTime} + return + } +} + +func (f *podLogFilter) cutoff(name, namespace string, gvk schema.GroupVersionKind) time.Time { + if f == nil { + return time.Time{} + } + f.mu.RLock() + defer f.mu.RUnlock() + return f.cutoffs[kdutil.ResourceID(name, namespace, gvk)].time +} diff --git a/pkg/kubedog/logs_test.go b/pkg/kubedog/logs_test.go new file mode 100644 index 00000000..15f907b2 --- /dev/null +++ b/pkg/kubedog/logs_test.go @@ -0,0 +1,206 @@ +package kubedog + +import ( + "sync/atomic" + "testing" + "time" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + "github.com/werf/kubedog/pkg/informer" + kdutil "github.com/werf/kubedog/pkg/trackers/dyntracker/util" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/apis/meta/v1/unstructured" + "k8s.io/apimachinery/pkg/runtime" + "k8s.io/apimachinery/pkg/runtime/schema" + "k8s.io/apimachinery/pkg/watch" + "k8s.io/client-go/dynamic/fake" + clienttesting "k8s.io/client-go/testing" + "k8s.io/client-go/tools/cache" +) + +func newLogFilterPod(name, namespace, uid, readyStatus string, readyTime time.Time) *unstructured.Unstructured { + pod := &unstructured.Unstructured{Object: map[string]any{ + "apiVersion": "v1", + "kind": "Pod", + "metadata": map[string]any{"name": name, "namespace": namespace, "uid": uid}, + "status": map[string]any{"phase": "Running"}, + }} + if readyStatus != "" { + condition := map[string]any{"type": "Ready", "status": readyStatus} + if !readyTime.IsZero() { + condition["lastTransitionTime"] = readyTime.Format(time.RFC3339Nano) + } + pod.Object["status"].(map[string]any)["conditions"] = []any{condition} + } + return pod +} + +func TestPodLogFilter_ReadyCondition(t *testing.T) { + readyTime := time.Date(2026, 10, 1, 12, 0, 0, 123, time.UTC) + for _, status := range []string{"", "False", "Unknown", "True"} { + t.Run("Ready="+status, func(t *testing.T) { + filter := newPodLogFilter() + filter.update(newLogFilterPod("app", "ns", "uid", status, readyTime)) + cutoff := filter.cutoff("app", "ns", watchdogPodGVK) + if status == "True" { + assert.Equal(t, readyTime, cutoff) + } else { + assert.True(t, cutoff.IsZero(), "Running alone does not mean that a pod is ready") + } + }) + } +} + +func TestPodLogFilter_FirstReadyTransitionIsRetained(t *testing.T) { + filter := newPodLogFilter() + firstReady := time.Date(2026, 10, 1, 12, 0, 0, 123, time.UTC) + filter.update(newLogFilterPod("app", "ns", "uid", "True", firstReady)) + filter.update(newLogFilterPod("app", "ns", "uid", "False", firstReady.Add(time.Minute))) + filter.update(newLogFilterPod("app", "ns", "uid", "True", firstReady.Add(2*time.Minute))) + assert.Equal(t, firstReady, filter.cutoff("app", "ns", watchdogPodGVK)) +} + +func TestPodLogFilter_JobPodsAreExempt(t *testing.T) { + filter := newPodLogFilter() + pod := newLogFilterPod("job-pod", "ns", "uid", "True", time.Now()) + pod.SetOwnerReferences([]metav1.OwnerReference{{APIVersion: "batch/v1", Kind: "Job", Name: "job"}}) + filter.update(pod) + assert.True(t, filter.cutoff("job-pod", "ns", watchdogPodGVK).IsZero()) +} + +func TestPodLogFilter_ReplacementPodStartsLoggingAgain(t *testing.T) { + filter := newPodLogFilter() + firstReady := time.Date(2026, 10, 1, 12, 0, 0, 123, time.UTC) + filter.update(newLogFilterPod("app", "ns", "old-uid", "True", firstReady)) + filter.update(newLogFilterPod("app", "ns", "new-uid", "False", time.Time{})) + assert.True(t, filter.cutoff("app", "ns", watchdogPodGVK).IsZero()) + + nextReady := firstReady.Add(time.Minute) + filter.update(newLogFilterPod("app", "ns", "new-uid", "True", nextReady)) + assert.Equal(t, nextReady, filter.cutoff("app", "ns", watchdogPodGVK)) +} + +func TestPodLogFilter_NamespacesAreIndependent(t *testing.T) { + filter := newPodLogFilter() + readyTime := time.Date(2026, 10, 1, 12, 0, 0, 123, time.UTC) + filter.update(newLogFilterPod("app", "ready-ns", "uid", "True", readyTime)) + filter.update(newLogFilterPod("app", "starting-ns", "uid", "False", time.Time{})) + assert.Equal(t, readyTime, filter.cutoff("app", "ready-ns", watchdogPodGVK)) + assert.True(t, filter.cutoff("app", "starting-ns", watchdogPodGVK).IsZero()) +} + +func TestPodLogFilter_PreservesReadinessSecond(t *testing.T) { + filter := newPodLogFilter() + readyTime := time.Date(2026, 10, 1, 12, 0, 0, 0, time.UTC) + filter.update(newLogFilterPod("app", "ns", "uid", "True", readyTime)) + assert.Equal(t, readyTime.Add(time.Second-time.Nanosecond), filter.cutoff("app", "ns", watchdogPodGVK)) +} + +func TestPodLogFilter_MissingOrInvalidTransitionTime(t *testing.T) { + for _, timestamp := range []string{"", "invalid"} { + t.Run(timestamp, func(t *testing.T) { + filter := newPodLogFilter() + pod := newLogFilterPod("app", "ns", "uid", "True", time.Time{}) + conditions := pod.Object["status"].(map[string]any)["conditions"].([]any) + conditions[0].(map[string]any)["lastTransitionTime"] = timestamp + before := time.Now() + filter.update(pod) + after := time.Now() + cutoff := filter.cutoff("app", "ns", watchdogPodGVK) + assert.False(t, cutoff.Before(before)) + assert.False(t, cutoff.After(after)) + }) + } +} + +func TestPodLogFilter_WatchesExistingPodsAndUpdatesUsingSharedInformer(t *testing.T) { + ctx := t.Context() + podGVR := schema.GroupVersionResource{Version: "v1", Resource: "pods"} + readyTime := time.Date(2026, 10, 1, 12, 0, 0, 123, time.UTC) + readyPod := newLogFilterPod("ready", "ns", "ready-uid", "True", readyTime) + startingPod := newLogFilterPod("starting", "ns", "starting-uid", "False", time.Time{}) + client := fake.NewSimpleDynamicClientWithCustomListKinds( + runtime.NewScheme(), map[schema.GroupVersionResource]string{podGVR: "PodList"}, readyPod, startingPod, + ) + var watches atomic.Int32 + client.PrependWatchReactor("pods", func(clienttesting.Action) (bool, watch.Interface, error) { + watches.Add(1) + return false, nil, nil + }) + factory := informer.NewConcurrentInformerFactory(ctx.Done(), make(chan error, 1), client, + informer.ConcurrentInformerFactoryOptions{}) + + // Start another consumer first, as kubedog does during resource tracking. + var shared *kdutil.Concurrent[*informer.Informer] + var err error + factory.RWTransaction(func(factory *informer.InformerFactory) { + shared, err = factory.ForNamespace(podGVR, "ns") + }) + require.NoError(t, err) + shared.RWTransaction(func(inf *informer.Informer) { + _, err = inf.AddEventHandler(cache.ResourceEventHandlerFuncs{}) + if err == nil { + inf.Run() + } + }) + require.NoError(t, err) + + filter := newPodLogFilter() + require.NoError(t, filter.watch(factory, []trackTarget{ + {kind: "deploy", namespace: "ns"}, + {kind: "sts", namespace: "ns"}, + {kind: "job", namespace: "jobs"}, + {kind: "pvc", namespace: "storage"}, + {kind: "canary", namespace: "canaries"}, + })) + require.Eventually(t, func() bool { + return filter.cutoff("ready", "ns", watchdogPodGVK).Equal(readyTime) + }, 5*time.Second, 10*time.Millisecond, "initial add must retain the pod's historical readiness time") + assert.True(t, filter.cutoff("starting", "ns", watchdogPodGVK).IsZero()) + + nextReady := readyTime.Add(time.Minute) + startingPod = newLogFilterPod("starting", "ns", "starting-uid", "True", nextReady) + _, err = client.Resource(podGVR).Namespace("ns").Update(ctx, startingPod, metav1.UpdateOptions{}) + require.NoError(t, err) + require.Eventually(t, func() bool { + return filter.cutoff("starting", "ns", watchdogPodGVK).Equal(nextReady) + }, 5*time.Second, 10*time.Millisecond, "updates must stop logs independently of the other pod") + assert.EqualValues(t, 1, watches.Load(), "consumers and repeated targets must share one pod watch") +} + +func TestCutoffTrackedKind(t *testing.T) { + for kind, want := range map[string]bool{ + "deploy": true, + "sts": true, + "ds": true, + "job": false, + "pvc": false, + "canary": false, + } { + t.Run(kind, func(t *testing.T) { + assert.Equal(t, want, cutoffTrackedKind(kind)) + }) + } +} + +func TestNewStartupLogFilter(t *testing.T) { + // Every combination below must yield a nil filter without touching the + // informer factory: startup-only cutoffs apply solely to full streaming. + nilFilterCases := []struct { + name string + options *TrackOptions + }{ + {name: "nil options", options: nil}, + {name: "no log streaming", options: NewTrackOptions().WithLogsUntilReady(true)}, + {name: "streaming without cutoff", options: NewTrackOptions().WithLogs(true)}, + {name: "failed-only mode", options: NewTrackOptions().WithFailedLogsOnly(true).WithLogsUntilReady(true)}, + } + for _, tt := range nilFilterCases { + t.Run(tt.name, func(t *testing.T) { + filter, err := newStartupLogFilter(tt.options, nil, nil) + require.NoError(t, err) + assert.Nil(t, filter) + }) + } +} diff --git a/pkg/kubedog/options.go b/pkg/kubedog/options.go index 63dbe0e8..fe9af98b 100644 --- a/pkg/kubedog/options.go +++ b/pkg/kubedog/options.go @@ -34,6 +34,9 @@ type TrackOptions struct { Timeout time.Duration // Logs enables emitting logs for every pod kubedog observes. Logs bool + // LogsUntilReady limits Logs to startup output through each pod's first + // Ready transition. Job logs are unaffected. Has no effect unless Logs is true. + LogsUntilReady bool // FailedLogsOnly enables capturing logs in the background and emitting // them only for pods that enter a failed state (CrashLoopBackOff, Error, // ImagePullBackOff, etc.). Has no effect when Logs is true. @@ -81,6 +84,12 @@ func (o *TrackOptions) WithLogsInterval(interval time.Duration) *TrackOptions { return o } +// WithLogsUntilReady limits pod log output to startup when log streaming is enabled. +func (o *TrackOptions) WithLogsUntilReady(v bool) *TrackOptions { + o.LogsUntilReady = v + 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 63a11d0b..b26c58b8 100644 --- a/pkg/kubedog/printer.go +++ b/pkg/kubedog/printer.go @@ -187,6 +187,7 @@ type progressPrinter struct { logStore *kdutil.Concurrent[*logstore.LogStore] skipLogs bool failedLogsOnly bool + podLogFilter *podLogFilter gates *gateStatuses skipped *skippedKeys useColor bool @@ -726,6 +727,7 @@ func (p *progressPrinter) flushLogs() { return } } + cutoff := p.podLogFilter.cutoff(rl.Name(), rl.Namespace(), rl.GroupVersionKind()) for source, lines := range rl.LogLines() { cursorKey := resourceKey + "|" + source start := p.lastCounts[cursorKey] @@ -733,6 +735,9 @@ func (p *progressPrinter) flushLogs() { continue } for _, ll := range lines[start:] { + if !cutoff.IsZero() && ll.Time.After(cutoff) { + continue + } pending = append(pending, entry{ resourceKey: resourceKey, displayLabel: displayLabel, diff --git a/pkg/kubedog/printer_test.go b/pkg/kubedog/printer_test.go index 368b70d2..21f0785a 100644 --- a/pkg/kubedog/printer_test.go +++ b/pkg/kubedog/printer_test.go @@ -15,6 +15,7 @@ import ( kdutil "github.com/werf/kubedog/pkg/trackers/dyntracker/util" "go.uber.org/zap" "go.uber.org/zap/zapcore" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/apimachinery/pkg/runtime/schema" ) @@ -511,6 +512,94 @@ func TestFlushLogs_HeaderDedupedAcrossFlushes(t *testing.T) { assert.Contains(t, out, "line 4") } +func TestFlushLogs_LogsUntilReady(t *testing.T) { + for _, kind := range []string{"Deployment", "StatefulSet", "DaemonSet"} { + t.Run(kind, func(t *testing.T) { + taskStore := kdutil.NewConcurrent(statestore.NewTaskStore()) + logStore := kdutil.NewConcurrent(logstore.NewLogStore()) + gvk := schema.GroupVersionKind{Group: "apps", Version: "v1", Kind: kind} + ts := statestore.NewReadinessTaskState("app", "ns", gvk, statestore.ReadinessTaskStateOptions{}) + // The real dyntracker leaves individual pod states unknown until + // the entire workload is ready, even when a replica is already Ready. + addPodChild(t, ts, "app-ready", "ns", statestore.ResourceStatusUnknown, "Running") + addPodChild(t, ts, "app-starting", "ns", statestore.ResourceStatusUnknown, "Running") + taskStore.RWTransaction(func(s *statestore.TaskStore) { + s.AddReadinessTaskState(kdutil.NewConcurrent(ts)) + }) + + readyTime := time.Unix(100, 500) + filter := newPodLogFilter() + filter.update(newLogFilterPod("app-ready", "ns", "ready-uid", "True", readyTime)) + filter.update(newLogFilterPod("app-starting", "ns", "starting-uid", "False", readyTime)) + logger, buf, mu := newBufferedLogger(t) + p := newProgressPrinter(logger, "", taskStore, logStore, false, false, newGateStatuses(), newSkippedKeys(), false, 0) + p.podLogFilter = filter + + readyLogs := kdutil.NewConcurrent(logstore.NewResourceLogs("app-ready", "ns", podGVK)) + startingLogs := kdutil.NewConcurrent(logstore.NewResourceLogs("app-starting", "ns", podGVK)) + logStore.RWTransaction(func(s *logstore.LogStore) { + s.AddResourceLogs(readyLogs) + s.AddResourceLogs(startingLogs) + }) + readyLogs.RWTransaction(func(rl *logstore.ResourceLogs) { + rl.AddLogLine("successful startup", "container/main", readyTime.Add(-time.Second)) + rl.AddLogLine("at readiness", "container/main", readyTime) + rl.AddLogLine("HTTP request", "container/main", readyTime.Add(time.Second)) + }) + startingLogs.RWTransaction(func(rl *logstore.ResourceLogs) { + rl.AddLogLine("still starting", "container/main", readyTime.Add(time.Second)) + }) + p.flushLogs() + + // Late pre-readiness chunks from another container still flush; + // repeated traffic from the ready replica stays silent. + readyLogs.RWTransaction(func(rl *logstore.ResourceLogs) { + rl.AddLogLine("late startup chunk", "container/sidecar", readyTime.Add(-time.Millisecond)) + rl.AddLogLine("background job", "container/sidecar", readyTime.Add(time.Second)) + rl.AddLogLine("another HTTP request", "container/main", readyTime.Add(2*time.Second)) + }) + p.flushLogs() + out := capturedOutput(buf, mu) + assert.Contains(t, out, "successful startup") + assert.Contains(t, out, "at readiness") + assert.Contains(t, out, "late startup chunk") + assert.Contains(t, out, "still starting") + assert.NotContains(t, out, "HTTP request") + assert.NotContains(t, out, "background job") + assert.Equal(t, 4, p.lastCounts["Pod/ns/app-ready|container/main"]) + assert.Equal(t, 2, p.lastCounts["Pod/ns/app-ready|container/sidecar"]) + p.flushLogs() + assert.Equal(t, out, capturedOutput(buf, mu), "startup lines must not be duplicated") + + // Filtering logs must not mutate readiness or hide a later failure. + assert.NotEqual(t, statestore.ReadinessTaskStatusReady, ts.Status()) + flipPodPhase(t, ts, "app-ready", "ns", "CrashLoopBackOff") + p.flushProgress() + assert.Contains(t, capturedOutput(buf, mu), "CrashLoopBackOff") + assert.Contains(t, p.collectFailedPodIDs(), kdutil.ResourceID("app-ready", "ns", podGVK)) + }) + } +} + +func TestFlushLogs_LogsUntilReady_JobLogs(t *testing.T) { + taskStore := kdutil.NewConcurrent(statestore.NewTaskStore()) + logStore := kdutil.NewConcurrent(logstore.NewLogStore()) + logger, buf, mu := newBufferedLogger(t) + p := newProgressPrinter(logger, "", taskStore, logStore, false, false, newGateStatuses(), newSkippedKeys(), false, 0) + p.podLogFilter = newPodLogFilter() + pod := newLogFilterPod("job-pod", "ns", "job-uid", "True", time.Unix(0, 0)) + pod.SetOwnerReferences([]metav1.OwnerReference{{Kind: "Job", Name: "job"}}) + p.podLogFilter.update(pod) + addPodLogs(t, logStore, "job-pod", "job started", "job still running") + p.flushLogs() + addPodLogs(t, logStore, "job-pod", "job completed") + p.flushLogs() + out := capturedOutput(buf, mu) + assert.Contains(t, out, "job started") + assert.Contains(t, out, "job still running") + assert.Contains(t, out, "job completed") +} + func TestHeaderDivider_FormatsBothBorders(t *testing.T) { assert.Equal(t, "========== title: ==========", HeaderDivider("title:")) } diff --git a/pkg/kubedog/tracker.go b/pkg/kubedog/tracker.go index 76585a40..a0749643 100644 --- a/pkg/kubedog/tracker.go +++ b/pkg/kubedog/tracker.go @@ -454,6 +454,11 @@ func (t *Tracker) TrackResources(ctx context.Context, resources []*resource.Reso gateStatuses := newGateStatuses() printer := newProgressPrinter(t.logger, t.releaseName, taskStore, logStore, ignoreLogs, failedLogsOnly, gateStatuses, t.skipped, t.trackOptions.Color, t.trackOptions.LogsInterval) + podLogFilter, err := newStartupLogFilter(t.trackOptions, informerFactory, targets) + if err != nil { + return err + } + printer.podLogFilter = podLogFilter printerDone := make(chan struct{}) // Spawn a parallel failure watchdog. It catches pods that genuinely diff --git a/pkg/kubedog/tracker_test.go b/pkg/kubedog/tracker_test.go index cee8b28f..1cad0d00 100644 --- a/pkg/kubedog/tracker_test.go +++ b/pkg/kubedog/tracker_test.go @@ -32,6 +32,7 @@ func TestNewTrackOptions(t *testing.T) { assert.NotNil(t, opts) assert.Equal(t, 5*time.Minute, opts.Timeout) assert.Equal(t, false, opts.Logs) + assert.False(t, opts.LogsUntilReady) assert.Equal(t, 10*time.Minute, opts.LogsSince) assert.Equal(t, 10*time.Second, opts.LogsInterval) } @@ -57,6 +58,12 @@ func TestTrackOptions_WithLogsInterval(t *testing.T) { assert.Equal(t, 3*time.Second, opts.LogsInterval) } +func TestTrackOptions_WithLogsUntilReady(t *testing.T) { + opts := NewTrackOptions().WithLogs(true).WithLogsUntilReady(true) + assert.True(t, opts.LogsUntilReady) + assert.True(t, opts.Logs) +} + func TestTrackOptions_Chaining(t *testing.T) { opts := NewTrackOptions() opts = opts. diff --git a/pkg/state/helmx.go b/pkg/state/helmx.go index 9929bf96..72ee26bc 100644 --- a/pkg/state/helmx.go +++ b/pkg/state/helmx.go @@ -313,6 +313,11 @@ func (st *HelmState) buildReleaseTracker(release *ReleaseSpec, opts *SyncOpts, u trackLogs = opts.TrackLogs } + trackLogsUntilReady := release.TrackLogsUntilReady != nil && *release.TrackLogsUntilReady + if release.TrackLogsUntilReady == nil && opts != nil { + trackLogsUntilReady = opts.TrackLogsUntilReady + } + trackFailedLogs := release.TrackFailedLogs != nil && *release.TrackFailedLogs if release.TrackFailedLogs == nil && opts != nil { trackFailedLogs = opts.TrackFailedLogs @@ -327,6 +332,7 @@ func (st *HelmState) buildReleaseTracker(release *ReleaseSpec, opts *SyncOpts, u trackOpts := kubedog.NewTrackOptions(). WithTimeout(timeout). WithLogs(trackLogs). + WithLogsUntilReady(trackLogsUntilReady). WithFailedLogsOnly(trackFailedLogs). WithFilterConfig(filterConfig). WithColor(useColor) diff --git a/pkg/state/helmx_test.go b/pkg/state/helmx_test.go index d4eb5382..02327f00 100644 --- a/pkg/state/helmx_test.go +++ b/pkg/state/helmx_test.go @@ -2,12 +2,15 @@ package state import ( "bytes" + "os" + "path/filepath" "testing" "github.com/stretchr/testify/require" "github.com/helmfile/helmfile/pkg/helmexec" "github.com/helmfile/helmfile/pkg/testutil" + "github.com/helmfile/helmfile/pkg/yaml" ) func TestAppendWaitForJobsFlags(t *testing.T) { @@ -797,3 +800,62 @@ func TestAppendServerSideFlagsForUpgrade(t *testing.T) { }) } } + +func TestBuildReleaseTracker_LogsUntilReady(t *testing.T) { + kubeconfig := filepath.Join(t.TempDir(), "kubeconfig") + require.NoError(t, os.WriteFile(kubeconfig, []byte(`apiVersion: v1 +kind: Config +clusters: +- name: test + cluster: + server: http://127.0.0.1:1 +contexts: +- name: test + context: + cluster: test +current-context: test +`), 0600)) + + tests := []struct { + name string + release string + opts *SyncOpts + want bool + }{ + { + name: "disabled by default", + }, + { + name: "cli enables startup logs", + opts: &SyncOpts{TrackLogs: true, TrackLogsUntilReady: true}, + want: true, + }, + { + name: "release enables startup logs", + release: "trackLogs: true\ntrackLogsUntilReady: true\n", + want: true, + }, + { + name: "release false overrides cli true", + release: "trackLogsUntilReady: false\n", + opts: &SyncOpts{TrackLogs: true, TrackLogsUntilReady: true}, + }, + { + name: "release true overrides cli false", + release: "trackLogsUntilReady: true\n", + opts: &SyncOpts{TrackLogs: true}, + want: true, + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + var release ReleaseSpec + require.NoError(t, yaml.Unmarshal([]byte(tt.release), &release)) + st := &HelmState{kubeconfig: kubeconfig} + _, opts, err := st.buildReleaseTracker(&release, tt.opts, false) + require.NoError(t, err) + require.Equal(t, tt.want, opts.LogsUntilReady) + }) + } +} diff --git a/pkg/state/state.go b/pkg/state/state.go index 0a94b331..2bdc9d5f 100644 --- a/pkg/state/state.go +++ b/pkg/state/state.go @@ -561,9 +561,11 @@ type ReleaseSpec struct { TrackTimeout *int `yaml:"trackTimeout,omitempty"` // TrackLogs enables log streaming with kubedog TrackLogs *bool `yaml:"trackLogs,omitempty"` + // TrackLogsUntilReady stops pod log streaming when each pod becomes ready. Requires TrackLogs. + TrackLogsUntilReady *bool `yaml:"trackLogsUntilReady,omitempty"` // TrackFailedLogs streams logs only for pods that enter a failed state // (CrashLoopBackOff, Error, etc.). Pods that succeed produce no output. - // Has no effect when TrackLogs is true (full streaming wins). + // Has no effect when TrackLogs is true (TrackLogs wins). TrackFailedLogs *bool `yaml:"trackFailedLogs,omitempty"` // HelmStuckGrace, when > 0, enables the safety-valve helm-killer for // kubedog tracking: if the cluster confirms every tracked resource has @@ -1082,6 +1084,7 @@ type SyncOpts struct { TrackMode string TrackTimeout int TrackLogs bool + TrackLogsUntilReady bool TrackFailedLogs bool TrackLogsInterval time.Duration HelmStuckGrace int diff --git a/pkg/state/temp_test.go b/pkg/state/temp_test.go index 178d3221..e4f35a27 100644 --- a/pkg/state/temp_test.go +++ b/pkg/state/temp_test.go @@ -42,39 +42,39 @@ func TestGenerateID(t *testing.T) { run(testcase{ subject: "baseline", release: ReleaseSpec{Name: "foo", Chart: "incubator/raw"}, - want: "foo-values-565b4d4448", + want: "foo-values-699f8dfb8f", }) run(testcase{ subject: "different bytes content", release: ReleaseSpec{Name: "foo", Chart: "incubator/raw"}, data: []byte(`{"k":"v"}`), - want: "foo-values-5bf759b9cf", + want: "foo-values-8484b97b54", }) run(testcase{ subject: "different map content", release: ReleaseSpec{Name: "foo", Chart: "incubator/raw"}, data: map[string]any{"k": "v"}, - want: "foo-values-cbd47bb9c", + want: "foo-values-6d74854796", }) run(testcase{ subject: "different chart", release: ReleaseSpec{Name: "foo", Chart: "stable/envoy"}, - want: "foo-values-84bcbdd8f9", + want: "foo-values-6c75d94b64", }) run(testcase{ subject: "different name", release: ReleaseSpec{Name: "bar", Chart: "incubator/raw"}, - want: "bar-values-757bdb6889", + want: "bar-values-fdc66db9d", }) run(testcase{ subject: "specific ns", release: ReleaseSpec{Name: "foo", Chart: "incubator/raw", Namespace: "myns"}, - want: "myns-foo-values-f4bcf8c77", + want: "myns-foo-values-74d9576cdb", }) for id, n := range ids {