add trackLogsUntilReady flag (#2814)

* add trackLogsUntilReady flag

Signed-off-by: pznamensky <kompastver@gmail.com>

* refactor(kubedog): flatten nested conditionals in trackLogsUntilReady filter

- podLogFilter.update: replace nested if-found/if-uid block with a single
  combined guard plus a no-op delete
- podLogFilter.watch: extract cutoffTrackedKind predicate and
  watchNamespacePods helper; use RWTransactionErr (kubedog's own idiom)
  instead of capturing err in plain RWTransaction closures, removing the
  'if err == nil { inf.Run() }' nesting
- TrackResources: extract newStartupLogFilter guard-clause helper so the
  printer wiring stays flat
- add table tests for the new helpers

No behavior change; race tests and golangci-lint pass.

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

---------

Signed-off-by: pznamensky <kompastver@gmail.com>
Signed-off-by: yxxhero <aiopsclub@163.com>
Co-authored-by: yxxhero <aiopsclub@163.com>
This commit is contained in:
pznamensky
2026-10-04 08:49:05 +08:00
committed by GitHub
co-authored by yxxhero
parent 000a6fa660
commit c25d16a438
22 changed files with 645 additions and 7 deletions
+17
View File
@@ -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:
+1
View File
@@ -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)")
+19
View File
@@ -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)
})
}
}
+1
View File
@@ -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)")
+27
View File
@@ -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:
+1
View File
@@ -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 |
+2
View File
@@ -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(),
+5
View File
@@ -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
}
+2
View File
@@ -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
+7
View File
@@ -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
+7
View File
@@ -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
+157
View File
@@ -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
}
+206
View File
@@ -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)
})
}
}
+9
View File
@@ -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
+5
View File
@@ -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,
+89
View File
@@ -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:"))
}
+5
View File
@@ -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
+7
View File
@@ -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.
+6
View File
@@ -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)
+62
View File
@@ -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)
})
}
}
+4 -1
View File
@@ -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
+6 -6
View File
@@ -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 {