diff --git a/pkg/model/khifile/v6/timeline_accumulator.go b/pkg/model/khifile/v6/timeline_accumulator.go index 63bf59ec1..d250e07b4 100644 --- a/pkg/model/khifile/v6/timeline_accumulator.go +++ b/pkg/model/khifile/v6/timeline_accumulator.go @@ -80,22 +80,6 @@ func (a *TimelineAccumulator) SetAlias(aliasPath, targetPath *TimelinePath) erro return a.registry.SetAlias(aliasPath, targetPath) } -// HasRevision reports whether the timeline at path has any accumulated revisions. -func (a *TimelineAccumulator) HasRevision(path *TimelinePath) bool { - if b, ok := a.registry.GetBuilderIfExists(path); ok { - return b.HasRevision() - } - return false -} - -// HasEvent reports whether the timeline at path has any accumulated events. -func (a *TimelineAccumulator) HasEvent(path *TimelinePath) bool { - if b, ok := a.registry.GetBuilderIfExists(path); ok { - return b.HasEvent() - } - return false -} - // NotifyItemsAdded increments the count of accumulated items and triggers a flush if the threshold is met. func (a *TimelineAccumulator) NotifyItemsAdded(count int) error { if count <= 0 { @@ -153,9 +137,3 @@ func (a *TimelineAccumulator) Flush() error { } return nil } - -// AddTestRevision adds a dummy revision to the timeline builder at path for testing purposes. -func (a *TimelineAccumulator) AddTestRevision(path *TimelinePath) { - b := a.GetBuilder(path) - b.AddRevision(pendingRevision{}) -} diff --git a/pkg/model/khifile/v6/timeline_builder.go b/pkg/model/khifile/v6/timeline_builder.go index 58761344c..22502d2c7 100644 --- a/pkg/model/khifile/v6/timeline_builder.go +++ b/pkg/model/khifile/v6/timeline_builder.go @@ -131,20 +131,6 @@ func (b *TimelineBuilder) HasItems() bool { return len(b.events) > 0 || len(b.revisions) > 0 } -// HasRevision returns true if the builder has accumulated any revisions. -func (b *TimelineBuilder) HasRevision() bool { - b.mu.Lock() - defer b.mu.Unlock() - return len(b.revisions) > 0 -} - -// HasEvent returns true if the builder has accumulated any events. -func (b *TimelineBuilder) HasEvent() bool { - b.mu.Lock() - defer b.mu.Unlock() - return len(b.events) > 0 -} - // HasEverHadItems returns true if the builder has ever accumulated events or revisions, // even if they have already been flushed from memory. func (b *TimelineBuilder) HasEverHadItems() bool { diff --git a/pkg/task/inspection/common/k8saudit/impl/inventory_timeline_creation_time.go b/pkg/task/inspection/common/k8saudit/impl/inventory_timeline_creation_time.go new file mode 100644 index 000000000..8981d4a59 --- /dev/null +++ b/pkg/task/inspection/common/k8saudit/impl/inventory_timeline_creation_time.go @@ -0,0 +1,170 @@ +// Copyright 2026 Google LLC +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +package k8saudit_impl + +import ( + "context" + "slices" + "time" + + "github.com/GoogleCloudPlatform/khi/pkg/core/inspection/progress" + inspectiontaskbase "github.com/GoogleCloudPlatform/khi/pkg/core/inspection/taskbase" + coretask "github.com/GoogleCloudPlatform/khi/pkg/core/task" + pb "github.com/GoogleCloudPlatform/khi/pkg/generated/khifile/v6" + "github.com/GoogleCloudPlatform/khi/pkg/task/inspection/common/k8saudit" + "github.com/GoogleCloudPlatform/khi/pkg/task/inspection/inspectioncore" +) + +// TimelineCreationTimeInventoryTask aggregates creation timestamps per timeline path discovered from audit logs. +var TimelineCreationTimeInventoryTask = inspectiontaskbase.NewInventoryTask( + k8saudit.TimelineCreationTimeInventoryTaskID, + k8saudit.TagTimelineCreationTimeDiscovery, + mergeTimelineCreationTimes, +) + +func mergeTimelineCreationTimes(results []k8saudit.TimelineCreationTimes) (k8saudit.TimelineCreationTimes, error) { + result := k8saudit.TimelineCreationTimes{} + for _, r := range results { + for path, times := range r { + result[path] = append(result[path], times...) + } + } + for path, times := range result { + result[path] = deduplicateAndSortTimes(times) + } + return result, nil +} + +func deduplicateAndSortTimes(times []time.Time) []time.Time { + if len(times) <= 1 { + return slices.Clone(times) + } + sorted := slices.Clone(times) + slices.SortFunc(sorted, func(a, b time.Time) int { + return a.Compare(b) + }) + return slices.CompactFunc(sorted, func(a, b time.Time) bool { + return a.Equal(b) + }) +} + +func extractLogCreationTime(l *k8saudit.ResourceManifestLog, verb *pb.Verb) (time.Time, bool) { + if l.ResourceBodyReader != nil { + if creationTime, found := GetCreationTimestamp(l.ResourceBodyReader); found { + return creationTime, true + } + } + if verb == k8saudit.VerbCreate { + return l.Log.Timestamp, true + } + return time.Time{}, false +} + +// ResourceTimelineCreationTimeDiscoveryTask extracts resource timeline creation timestamps from audit logs. +var ResourceTimelineCreationTimeDiscoveryTask = inspectiontaskbase.NewInspectionTask( + k8saudit.ResourceTimelineCreationTimeDiscoveryTaskID, + []coretask.Dependency{ + k8saudit.ManifestGeneratorTaskID.Ref(), + k8saudit.K8sAuditLogExtractorRef.Ref(coretask.FromActiveGraph), + }, + func(ctx context.Context, taskMode inspectioncore.InspectionTaskModeType) (k8saudit.TimelineCreationTimes, error) { + if taskMode == inspectioncore.TaskModeDryRun { + return k8saudit.TimelineCreationTimes{}, nil + } + result := k8saudit.TimelineCreationTimes{} + resourceLogs := coretask.GetTaskResult(ctx, k8saudit.ManifestGeneratorTaskID.Ref()) + tracker := progress.NewTracker(ctx, len(resourceLogs), progress.WithUnit("groups")) + defer tracker.Done() + for _, group := range resourceLogs { + if group.Resource.Type() == k8saudit.Namespace { + tracker.Inc() + continue + } + for _, l := range group.Logs { + k8sFieldSet, err := k8saudit.ExtractK8sAuditLog(ctx, l.Log.NodeReader) + if err != nil || k8sFieldSet == nil || k8sFieldSet.IsDryRun { + continue + } + creationTime, hasCreationTime := extractLogCreationTime(l, k8sFieldSet.Verb) + if !hasCreationTime { + continue + } + targetPath := MustResolveTimelinePath(ctx, k8sFieldSet.ClusterName, group.Resource) + result[targetPath] = append(result[targetPath], creationTime) + } + tracker.Inc() + } + for path, times := range result { + result[path] = deduplicateAndSortTimes(times) + } + return result, nil + }, + coretask.ProvidesTag(k8saudit.TagTimelineCreationTimeDiscovery), + coretask.WithFeatureGate(k8saudit.K8sAuditLogParserTailRef), + progress.WithTitle("Discover resource timeline creation times"), + coretask.WithTaskDescription("Extracts resource timeline creation timestamps from Kubernetes audit logs."), +) + +// PodPhaseTimelineCreationTimeDiscoveryTask extracts Pod phase timeline creation timestamps from audit logs. +var PodPhaseTimelineCreationTimeDiscoveryTask = inspectiontaskbase.NewInspectionTask( + k8saudit.PodPhaseTimelineCreationTimeDiscoveryTaskID, + []coretask.Dependency{ + k8saudit.ManifestGeneratorTaskID.Ref(), + k8saudit.K8sAuditLogExtractorRef.Ref(coretask.FromActiveGraph), + }, + func(ctx context.Context, taskMode inspectioncore.InspectionTaskModeType) (k8saudit.TimelineCreationTimes, error) { + if taskMode == inspectioncore.TaskModeDryRun { + return k8saudit.TimelineCreationTimes{}, nil + } + result := k8saudit.TimelineCreationTimes{} + resourceLogs := coretask.GetTaskResult(ctx, k8saudit.ManifestGeneratorTaskID.Ref()) + tracker := progress.NewTracker(ctx, len(resourceLogs), progress.WithUnit("groups")) + defer tracker.Done() + for _, group := range resourceLogs { + if group.Resource.Type() != k8saudit.Resource || group.Resource.APIVersion != "core/v1" || group.Resource.Kind != "pod" { + tracker.Inc() + continue + } + for _, l := range group.Logs { + if l.ResourceBodyReader == nil { + continue + } + k8sFieldSet, err := k8saudit.ExtractK8sAuditLog(ctx, l.Log.NodeReader) + if err != nil || k8sFieldSet == nil || k8sFieldSet.IsDryRun { + continue + } + creationTime, hasCreationTime := extractLogCreationTime(l, k8sFieldSet.Verb) + if !hasCreationTime { + continue + } + uid, foundUID := GetUID(l.ResourceBodyReader) + nodeName, foundNode := GetNodeNameOfPod(l.ResourceBodyReader) + if foundUID && uid != "" && foundNode && nodeName != "" { + podPhasePath := MustPodPhaseTimelinePath(ctx, k8sFieldSet.ClusterName, nodeName, group.Resource.Namespace, group.Resource.Name, uid) + result[podPhasePath] = append(result[podPhasePath], creationTime) + } + } + tracker.Inc() + } + for path, times := range result { + result[path] = deduplicateAndSortTimes(times) + } + return result, nil + }, + coretask.ProvidesTag(k8saudit.TagTimelineCreationTimeDiscovery), + coretask.WithFeatureGate(k8saudit.K8sAuditLogParserTailRef), + progress.WithTitle("Discover Pod phase timeline creation times"), + coretask.WithTaskDescription("Extracts Pod phase timeline creation timestamps from Kubernetes audit logs."), +) diff --git a/pkg/task/inspection/common/k8saudit/impl/inventory_timeline_creation_time_test.go b/pkg/task/inspection/common/k8saudit/impl/inventory_timeline_creation_time_test.go new file mode 100644 index 000000000..893fee307 --- /dev/null +++ b/pkg/task/inspection/common/k8saudit/impl/inventory_timeline_creation_time_test.go @@ -0,0 +1,646 @@ +// Copyright 2026 Google LLC +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +package k8saudit_impl + +import ( + "context" + "testing" + "time" + + "github.com/GoogleCloudPlatform/khi/pkg/common/structured" + inspectiontest "github.com/GoogleCloudPlatform/khi/pkg/core/inspection/test" + tasktest "github.com/GoogleCloudPlatform/khi/pkg/core/task/test" + pb "github.com/GoogleCloudPlatform/khi/pkg/generated/khifile/v6" + khifilev6 "github.com/GoogleCloudPlatform/khi/pkg/model/khifile/v6" + "github.com/GoogleCloudPlatform/khi/pkg/task/inspection/common/k8saudit" + "github.com/GoogleCloudPlatform/khi/pkg/task/inspection/inspectioncore" + "github.com/GoogleCloudPlatform/khi/pkg/testutil/testlog" + "github.com/google/go-cmp/cmp" +) + +func TestMergeTimelineCreationTimes(t *testing.T) { + ctx := inspectiontest.WithDefaultTestInspectionTaskContext(t.Context()) + p1 := MustResolveTimelinePath(ctx, "c", &k8saudit.ResourceIdentity{APIVersion: "core/v1", Kind: "pod", Namespace: "default", Name: "p1"}) + p2 := MustResolveTimelinePath(ctx, "c", &k8saudit.ResourceIdentity{APIVersion: "core/v1", Kind: "pod", Namespace: "default", Name: "p2"}) + p3 := MustResolveTimelinePath(ctx, "c", &k8saudit.ResourceIdentity{APIVersion: "core/v1", Kind: "pod", Namespace: "default", Name: "p3"}) + + t1 := time.Date(2026, 5, 26, 10, 0, 0, 0, time.UTC) + t2 := time.Date(2026, 5, 26, 11, 0, 0, 0, time.UTC) + t3 := time.Date(2026, 5, 26, 12, 0, 0, 0, time.UTC) + + testCases := []struct { + name string + inputs []k8saudit.TimelineCreationTimes + want k8saudit.TimelineCreationTimes + }{ + { + name: "merge disjoint sets", + inputs: []k8saudit.TimelineCreationTimes{ + {p1: []time.Time{t1}}, + {p2: []time.Time{t2}, p3: []time.Time{t3}}, + }, + want: k8saudit.TimelineCreationTimes{ + p1: []time.Time{t1}, + p2: []time.Time{t2}, + p3: []time.Time{t3}, + }, + }, + { + name: "merge overlapping sets with duplicate and out-of-order timestamps", + inputs: []k8saudit.TimelineCreationTimes{ + {p1: []time.Time{t2, t1}}, + {p1: []time.Time{t1, t3}, p2: []time.Time{t1}}, + }, + want: k8saudit.TimelineCreationTimes{ + p1: []time.Time{t1, t2, t3}, + p2: []time.Time{t1}, + }, + }, + { + name: "empty inputs", + inputs: []k8saudit.TimelineCreationTimes{}, + want: k8saudit.TimelineCreationTimes{}, + }, + } + + for _, tc := range testCases { + t.Run(tc.name, func(t *testing.T) { + got, err := mergeTimelineCreationTimes(tc.inputs) + if err != nil { + t.Fatalf("mergeTimelineCreationTimes() unexpected error: %v", err) + } + if diff := cmp.Diff(tc.want, got, cmp.Comparer(func(a, b *khifilev6.TimelinePath) bool { return a == b })); diff != "" { + t.Errorf("mergeTimelineCreationTimes() mismatch (-want +got):\n%s", diff) + } + }) + } +} + +func TestResourceTimelineCreationTimeDiscoveryTask(t *testing.T) { + testTime := time.Date(2026, 5, 26, 12, 0, 0, 0, time.UTC) + + type logSpec struct { + manifest string + verb *pb.Verb + logTime time.Time + isDryRun bool + } + + testCases := []struct { + name string + resource *k8saudit.ResourceIdentity + logs []logSpec + clusterName string + wantFunc func(ctx context.Context) k8saudit.TimelineCreationTimes + }{ + { + name: "standard pod resource with creationTimestamp", + resource: &k8saudit.ResourceIdentity{ + APIVersion: "core/v1", + Kind: "pod", + Namespace: "default", + Name: "test-pod", + }, + logs: []logSpec{ + { + manifest: `apiVersion: v1 +kind: Pod +metadata: + name: test-pod + namespace: default + uid: pod-uid-123 + creationTimestamp: "2026-05-26T10:00:00Z" +spec: + nodeName: worker-node-1 +`, + }, + }, + clusterName: "test-cluster", + wantFunc: func(ctx context.Context) k8saudit.TimelineCreationTimes { + podPath := MustResolveTimelinePath(ctx, "test-cluster", &k8saudit.ResourceIdentity{ + APIVersion: "core/v1", + Kind: "pod", + Namespace: "default", + Name: "test-pod", + }) + return k8saudit.TimelineCreationTimes{ + podPath: []time.Time{time.Date(2026, 5, 26, 10, 0, 0, 0, time.UTC)}, + } + }, + }, + { + name: "recreated pod on the same timeline path with multiple creationTimestamps and duplicate logs", + resource: &k8saudit.ResourceIdentity{ + APIVersion: "core/v1", + Kind: "pod", + Namespace: "default", + Name: "test-pod", + }, + logs: []logSpec{ + { + manifest: `apiVersion: v1 +kind: Pod +metadata: + name: test-pod + namespace: default + uid: pod-uid-2 + creationTimestamp: "2026-05-26T11:00:00Z" +spec: {} +`, + }, + { + manifest: `apiVersion: v1 +kind: Pod +metadata: + name: test-pod + namespace: default + uid: pod-uid-1 + creationTimestamp: "2026-05-26T10:00:00Z" +spec: {} +`, + }, + { + manifest: `apiVersion: v1 +kind: Pod +metadata: + name: test-pod + namespace: default + uid: pod-uid-1 + creationTimestamp: "2026-05-26T10:00:00Z" +spec: {} +`, + }, + }, + clusterName: "test-cluster", + wantFunc: func(ctx context.Context) k8saudit.TimelineCreationTimes { + podPath := MustResolveTimelinePath(ctx, "test-cluster", &k8saudit.ResourceIdentity{ + APIVersion: "core/v1", + Kind: "pod", + Namespace: "default", + Name: "test-pod", + }) + return k8saudit.TimelineCreationTimes{ + podPath: []time.Time{ + time.Date(2026, 5, 26, 10, 0, 0, 0, time.UTC), + time.Date(2026, 5, 26, 11, 0, 0, 0, time.UTC), + }, + } + }, + }, + { + name: "pod log without creationTimestamp and non-create verb is excluded", + resource: &k8saudit.ResourceIdentity{ + APIVersion: "core/v1", + Kind: "pod", + Namespace: "default", + Name: "test-pod", + }, + logs: []logSpec{ + { + manifest: `apiVersion: v1 +kind: Pod +metadata: + name: test-pod + namespace: default +`, + verb: k8saudit.VerbDelete, + }, + }, + clusterName: "test-cluster", + wantFunc: func(ctx context.Context) k8saudit.TimelineCreationTimes { + return k8saudit.TimelineCreationTimes{} + }, + }, + { + name: "pod binding subresource with VerbCreate falls back to log timestamp when creationTimestamp is absent", + resource: &k8saudit.ResourceIdentity{ + APIVersion: "core/v1", + Kind: "pod", + Namespace: "default", + Name: "test-pod", + SubresourceName: "binding", + }, + logs: []logSpec{ + { + manifest: `apiVersion: v1 +kind: Binding +metadata: + name: test-pod + namespace: default +target: + kind: Node + name: worker-node-1 +`, + verb: k8saudit.VerbCreate, + logTime: testTime, + }, + }, + clusterName: "test-cluster", + wantFunc: func(ctx context.Context) k8saudit.TimelineCreationTimes { + bindingPath := MustResolveTimelinePath(ctx, "test-cluster", &k8saudit.ResourceIdentity{ + APIVersion: "core/v1", + Kind: "pod", + Namespace: "default", + Name: "test-pod", + SubresourceName: "binding", + }) + return k8saudit.TimelineCreationTimes{ + bindingPath: []time.Time{testTime}, + } + }, + }, + { + name: "dry run log is excluded", + resource: &k8saudit.ResourceIdentity{ + APIVersion: "core/v1", + Kind: "pod", + Namespace: "default", + Name: "test-pod", + }, + logs: []logSpec{ + { + manifest: `apiVersion: v1 +kind: Pod +metadata: + name: test-pod + namespace: default + creationTimestamp: "2026-05-26T10:00:00Z" +`, + isDryRun: true, + }, + }, + clusterName: "test-cluster", + wantFunc: func(ctx context.Context) k8saudit.TimelineCreationTimes { + return k8saudit.TimelineCreationTimes{} + }, + }, + { + name: "namespace resource type is skipped", + resource: &k8saudit.ResourceIdentity{ + APIVersion: "core/v1", + Kind: "pod", + Namespace: "default", + Name: "", + }, + logs: []logSpec{ + { + manifest: `apiVersion: v1 +kind: Namespace +metadata: + name: kube-system + creationTimestamp: "2026-05-26T10:00:00Z" +`, + }, + }, + clusterName: "test-cluster", + wantFunc: func(ctx context.Context) k8saudit.TimelineCreationTimes { + return k8saudit.TimelineCreationTimes{} + }, + }, + } + + for _, tc := range testCases { + t.Run(tc.name, func(t *testing.T) { + ctx := inspectiontest.WithDefaultTestInspectionTaskContext(t.Context()) + + var manifestLogs []*k8saudit.ResourceManifestLog + for _, ls := range tc.logs { + logTime := ls.logTime + if logTime.IsZero() { + logTime = testTime + } + fieldSet := &k8saudit.K8sAuditLogFieldSet{ + ClusterName: tc.clusterName, + IsDryRun: ls.isDryRun, + Verb: ls.verb, + } + l := testlog.NewMockLog(logTime, fieldSet) + + var reader *structured.NodeReader + if ls.manifest != "" { + yamlNode, err := structured.FromYAML(ls.manifest) + if err != nil { + t.Fatalf("failed to parse yaml: %v", err) + } + reader = structured.NewNodeReader(yamlNode) + } + manifestLogs = append(manifestLogs, &k8saudit.ResourceManifestLog{ + Log: l, + ResourceBodyReader: reader, + }) + } + + input := k8saudit.ResourceManifestLogGroupMap{ + "test": &k8saudit.ResourceManifestLogGroup{ + Resource: tc.resource, + Logs: manifestLogs, + }, + } + + mockExtractor := k8saudit.K8sAuditLogExtractor(func(reader *structured.NodeReader) (*k8saudit.K8sAuditLogFieldSet, error) { + if mock, ok := structured.GetMock[*k8saudit.K8sAuditLogFieldSet](reader); ok { + return mock, nil + } + return nil, nil + }) + + got, _, err := inspectiontest.RunInspectionTask(ctx, ResourceTimelineCreationTimeDiscoveryTask, inspectioncore.TaskModeRun, map[string]any{}, + tasktest.NewTaskDependencyValuePair(k8saudit.ManifestGeneratorTaskID.Ref(), input), + tasktest.NewTaskDependencyValuePair(k8saudit.K8sAuditLogExtractorRef, mockExtractor), + ) + if err != nil { + t.Fatalf("RunInspectionTask failed: %v", err) + } + + want := tc.wantFunc(ctx) + if diff := cmp.Diff(want, got, cmp.Comparer(func(a, b *khifilev6.TimelinePath) bool { return a == b })); diff != "" { + t.Errorf("ResourceTimelineCreationTimeDiscoveryTask mismatch (-want +got):\n%s", diff) + } + }) + } +} + +func TestPodPhaseTimelineCreationTimeDiscoveryTask(t *testing.T) { + testTime := time.Date(2026, 5, 26, 12, 0, 0, 0, time.UTC) + + type logSpec struct { + manifest string + verb *pb.Verb + logTime time.Time + isDryRun bool + } + + testCases := []struct { + name string + resource *k8saudit.ResourceIdentity + logs []logSpec + clusterName string + wantFunc func(ctx context.Context) k8saudit.TimelineCreationTimes + }{ + { + name: "scheduled pod with creationTimestamp, uid, and spec.nodeName", + resource: &k8saudit.ResourceIdentity{ + APIVersion: "core/v1", + Kind: "pod", + Namespace: "default", + Name: "test-pod", + }, + logs: []logSpec{ + { + manifest: `apiVersion: v1 +kind: Pod +metadata: + name: test-pod + namespace: default + uid: pod-uid-123 + creationTimestamp: "2026-05-26T10:00:00Z" +spec: + nodeName: worker-node-1 +`, + }, + }, + clusterName: "test-cluster", + wantFunc: func(ctx context.Context) k8saudit.TimelineCreationTimes { + podPhasePath := MustPodPhaseTimelinePath(ctx, "test-cluster", "worker-node-1", "default", "test-pod", "pod-uid-123") + return k8saudit.TimelineCreationTimes{ + podPhasePath: []time.Time{time.Date(2026, 5, 26, 10, 0, 0, 0, time.UTC)}, + } + }, + }, + { + name: "multiple logs with duplicate and multiple timestamps deduplicated and sorted", + resource: &k8saudit.ResourceIdentity{ + APIVersion: "core/v1", + Kind: "pod", + Namespace: "default", + Name: "test-pod", + }, + logs: []logSpec{ + { + manifest: `apiVersion: v1 +kind: Pod +metadata: + name: test-pod + namespace: default + uid: pod-uid-123 + creationTimestamp: "2026-05-26T11:00:00Z" +spec: + nodeName: worker-node-1 +`, + }, + { + manifest: `apiVersion: v1 +kind: Pod +metadata: + name: test-pod + namespace: default + uid: pod-uid-123 + creationTimestamp: "2026-05-26T10:00:00Z" +spec: + nodeName: worker-node-1 +`, + }, + { + manifest: `apiVersion: v1 +kind: Pod +metadata: + name: test-pod + namespace: default + uid: pod-uid-123 + creationTimestamp: "2026-05-26T10:00:00Z" +spec: + nodeName: worker-node-1 +`, + }, + }, + clusterName: "test-cluster", + wantFunc: func(ctx context.Context) k8saudit.TimelineCreationTimes { + podPhasePath := MustPodPhaseTimelinePath(ctx, "test-cluster", "worker-node-1", "default", "test-pod", "pod-uid-123") + return k8saudit.TimelineCreationTimes{ + podPhasePath: []time.Time{ + time.Date(2026, 5, 26, 10, 0, 0, 0, time.UTC), + time.Date(2026, 5, 26, 11, 0, 0, 0, time.UTC), + }, + } + }, + }, + { + name: "unscheduled pod without spec.nodeName is excluded", + resource: &k8saudit.ResourceIdentity{ + APIVersion: "core/v1", + Kind: "pod", + Namespace: "default", + Name: "test-pod", + }, + logs: []logSpec{ + { + manifest: `apiVersion: v1 +kind: Pod +metadata: + name: test-pod + namespace: default + uid: pod-uid-123 + creationTimestamp: "2026-05-26T10:00:00Z" +spec: {} +`, + }, + }, + clusterName: "test-cluster", + wantFunc: func(ctx context.Context) k8saudit.TimelineCreationTimes { + return k8saudit.TimelineCreationTimes{} + }, + }, + { + name: "pod without creationTimestamp and non-create verb is excluded", + resource: &k8saudit.ResourceIdentity{ + APIVersion: "core/v1", + Kind: "pod", + Namespace: "default", + Name: "test-pod", + }, + logs: []logSpec{ + { + manifest: `apiVersion: v1 +kind: Pod +metadata: + name: test-pod + namespace: default + uid: pod-uid-123 +spec: + nodeName: worker-node-1 +`, + verb: k8saudit.VerbDelete, + }, + }, + clusterName: "test-cluster", + wantFunc: func(ctx context.Context) k8saudit.TimelineCreationTimes { + return k8saudit.TimelineCreationTimes{} + }, + }, + { + name: "non-pod resource skipped", + resource: &k8saudit.ResourceIdentity{ + APIVersion: "core/v1", + Kind: "configmap", + Namespace: "default", + Name: "test-cm", + }, + logs: []logSpec{ + { + manifest: `apiVersion: v1 +kind: ConfigMap +metadata: + name: test-cm + namespace: default + uid: cm-uid-123 + creationTimestamp: "2026-05-26T10:00:00Z" +`, + }, + }, + clusterName: "test-cluster", + wantFunc: func(ctx context.Context) k8saudit.TimelineCreationTimes { + return k8saudit.TimelineCreationTimes{} + }, + }, + { + name: "dry run log is excluded", + resource: &k8saudit.ResourceIdentity{ + APIVersion: "core/v1", + Kind: "pod", + Namespace: "default", + Name: "test-pod", + }, + logs: []logSpec{ + { + manifest: `apiVersion: v1 +kind: Pod +metadata: + name: test-pod + namespace: default + uid: pod-uid-123 + creationTimestamp: "2026-05-26T10:00:00Z" +spec: + nodeName: worker-node-1 +`, + isDryRun: true, + }, + }, + clusterName: "test-cluster", + wantFunc: func(ctx context.Context) k8saudit.TimelineCreationTimes { + return k8saudit.TimelineCreationTimes{} + }, + }, + } + + for _, tc := range testCases { + t.Run(tc.name, func(t *testing.T) { + ctx := inspectiontest.WithDefaultTestInspectionTaskContext(t.Context()) + + var manifestLogs []*k8saudit.ResourceManifestLog + for _, ls := range tc.logs { + logTime := ls.logTime + if logTime.IsZero() { + logTime = testTime + } + fieldSet := &k8saudit.K8sAuditLogFieldSet{ + ClusterName: tc.clusterName, + IsDryRun: ls.isDryRun, + Verb: ls.verb, + } + l := testlog.NewMockLog(logTime, fieldSet) + + var reader *structured.NodeReader + if ls.manifest != "" { + yamlNode, err := structured.FromYAML(ls.manifest) + if err != nil { + t.Fatalf("failed to parse yaml: %v", err) + } + reader = structured.NewNodeReader(yamlNode) + } + manifestLogs = append(manifestLogs, &k8saudit.ResourceManifestLog{ + Log: l, + ResourceBodyReader: reader, + }) + } + + input := k8saudit.ResourceManifestLogGroupMap{ + "test": &k8saudit.ResourceManifestLogGroup{ + Resource: tc.resource, + Logs: manifestLogs, + }, + } + + mockExtractor := k8saudit.K8sAuditLogExtractor(func(reader *structured.NodeReader) (*k8saudit.K8sAuditLogFieldSet, error) { + if mock, ok := structured.GetMock[*k8saudit.K8sAuditLogFieldSet](reader); ok { + return mock, nil + } + return nil, nil + }) + + got, _, err := inspectiontest.RunInspectionTask(ctx, PodPhaseTimelineCreationTimeDiscoveryTask, inspectioncore.TaskModeRun, map[string]any{}, + tasktest.NewTaskDependencyValuePair(k8saudit.ManifestGeneratorTaskID.Ref(), input), + tasktest.NewTaskDependencyValuePair(k8saudit.K8sAuditLogExtractorRef, mockExtractor), + ) + if err != nil { + t.Fatalf("RunInspectionTask failed: %v", err) + } + + want := tc.wantFunc(ctx) + if diff := cmp.Diff(want, got, cmp.Comparer(func(a, b *khifilev6.TimelinePath) bool { return a == b })); diff != "" { + t.Errorf("PodPhaseTimelineCreationTimeDiscoveryTask mismatch (-want +got):\n%s", diff) + } + }) + } +} diff --git a/pkg/task/inspection/common/k8saudit/impl/mapper_error_test.go b/pkg/task/inspection/common/k8saudit/impl/mapper_error_test.go index f8c314dd2..52c7c84c9 100644 --- a/pkg/task/inspection/common/k8saudit/impl/mapper_error_test.go +++ b/pkg/task/inspection/common/k8saudit/impl/mapper_error_test.go @@ -175,7 +175,8 @@ func TestNonSuccessLogLogToTimelineMapperTask(t *testing.T) { ns := builder.TimelineAccumulator.GetPath(kind, khifilev6.PathSegment{Name: "default", Type: inspectioncore.TimelineTypeNamespace}) wantPath := builder.TimelineAccumulator.GetPath(ns, khifilev6.PathSegment{Name: "pod-1", Type: inspectioncore.TimelineTypeResource}) - if !builder.TimelineAccumulator.HasEvent(wantPath) { + protoItems := builder.TimelineAccumulator.GetBuilder(wantPath).ToProto() + if protoItems == nil || len(protoItems.GetEvents()) == 0 { t.Errorf("expected timeline %v to have events, but none found", wantPath) } }) diff --git a/pkg/task/inspection/common/k8saudit/impl/registration.go b/pkg/task/inspection/common/k8saudit/impl/registration.go index 4b1bfb44f..524e39833 100644 --- a/pkg/task/inspection/common/k8saudit/impl/registration.go +++ b/pkg/task/inspection/common/k8saudit/impl/registration.go @@ -57,5 +57,8 @@ func Register(registry coreinspection.InspectionTaskRegistry) error { ContainerIDPatternFinderTask, IPLeaseHistoryInventoryTask, IPLeaseHistoryDiscoveryTask, + TimelineCreationTimeInventoryTask, + ResourceTimelineCreationTimeDiscoveryTask, + PodPhaseTimelineCreationTimeDiscoveryTask, ) } diff --git a/pkg/task/inspection/common/k8saudit/inventory.go b/pkg/task/inspection/common/k8saudit/inventory.go index f7ecf8a9d..5cb9b8259 100644 --- a/pkg/task/inspection/common/k8saudit/inventory.go +++ b/pkg/task/inspection/common/k8saudit/inventory.go @@ -15,9 +15,12 @@ package k8saudit import ( + "time" + coretask "github.com/GoogleCloudPlatform/khi/pkg/core/task" "github.com/GoogleCloudPlatform/khi/pkg/core/task/taskid" "github.com/GoogleCloudPlatform/khi/pkg/model/history/resourceinfo/resourcelease" + khifilev6 "github.com/GoogleCloudPlatform/khi/pkg/model/khifile/v6" ) // TagNodeNameDiscovery is the tag for discovery tasks producing node names. @@ -45,3 +48,12 @@ type IPLeaseHistory = *resourcelease.ResourceLeaseHistory[*ResourceIdentity] var TagIPLeaseHistoryDiscovery = coretask.NewTag[IPLeaseHistory]("khi.google.com/inspection/commonlogk8saudit/iplease") var IPLeaseHistoryInventoryTaskID = taskid.NewDefaultImplementationID[IPLeaseHistory](TaskIDPrefix + "ip-lease-history-inventory") + +// TimelineCreationTimes maps a timeline path to its observed creation timestamps in chronological order. +type TimelineCreationTimes = map[*khifilev6.TimelinePath][]time.Time + +// TagTimelineCreationTimeDiscovery is the tag for discovery tasks producing creation timestamps per timeline path. +var TagTimelineCreationTimeDiscovery = coretask.NewTag[TimelineCreationTimes]("khi.google.com/inspection/commonlogk8saudit/timelinecreationtime") + +// TimelineCreationTimeInventoryTaskID is the task ID for the inventory task aggregating creation timestamps per timeline path. +var TimelineCreationTimeInventoryTaskID = taskid.NewDefaultImplementationID[TimelineCreationTimes](TaskIDPrefix + "timeline-creation-time-inventory") diff --git a/pkg/task/inspection/common/k8saudit/taskid.go b/pkg/task/inspection/common/k8saudit/taskid.go index 2b5019859..771ab6490 100644 --- a/pkg/task/inspection/common/k8saudit/taskid.go +++ b/pkg/task/inspection/common/k8saudit/taskid.go @@ -114,3 +114,9 @@ var ContainerIDPatternFinderTaskID = taskid.NewDefaultImplementationID[patternfi // IPLeaseHistoryDiscoveryTaskID is the task ID for extracting IP lease history from audit logs. var IPLeaseHistoryDiscoveryTaskID = taskid.NewDefaultImplementationID[IPLeaseHistory](TaskIDPrefix + "ip-lease-history-discovery") + +// ResourceTimelineCreationTimeDiscoveryTaskID is the task ID for extracting resource timeline creation timestamps from audit logs. +var ResourceTimelineCreationTimeDiscoveryTaskID = taskid.NewDefaultImplementationID[TimelineCreationTimes](TaskIDPrefix + "resource-timeline-creation-time-discovery") + +// PodPhaseTimelineCreationTimeDiscoveryTaskID is the task ID for extracting Pod phase timeline creation timestamps from audit logs. +var PodPhaseTimelineCreationTimeDiscoveryTaskID = taskid.NewDefaultImplementationID[TimelineCreationTimes](TaskIDPrefix + "pod-phase-timeline-creation-time-discovery") diff --git a/pkg/task/inspection/googlecloud/k8scontainer/impl/mapper.go b/pkg/task/inspection/googlecloud/k8scontainer/impl/mapper.go index 117d95f22..7b6e17cf8 100644 --- a/pkg/task/inspection/googlecloud/k8scontainer/impl/mapper.go +++ b/pkg/task/inspection/googlecloud/k8scontainer/impl/mapper.go @@ -168,8 +168,7 @@ var pathMetadataUID = structured.CompileFieldPath("metadata.uid") func (m *containerLogPodPhaseTimelineMapper) Dependencies() []coretask.Dependency { return []coretask.Dependency{ k8scontainer.ClusterIdentityTaskID.Ref(), - k8saudit.ResourceRevisionLogToTimelineMapperTaskID.Ref(), - k8saudit.PodPhaseLogToTimelineMapperTaskID.Ref(), + k8saudit.TimelineCreationTimeInventoryTaskID.Ref(), k8saudit.InitialResourceStateProviderRef, } } @@ -205,13 +204,7 @@ func (m *containerLogPodPhaseTimelineMapper) ProcessLogByGroup(ctx context.Conte clusterName = clusterIdentity.ClusterName } - // Construct paths for Pod and its binding - cluster := k8saudit.MustK8sClusterTimeline(ctx, clusterName) - api := k8saudit.MustK8sAPIVersionTimeline(ctx, cluster, "core/v1") - kind := k8saudit.MustK8sKindTimeline(ctx, api, "pod") - ns := k8saudit.MustK8sNamespaceTimeline(ctx, kind, containerFields.Namespace) - podPath := k8saudit.MustK8sNamespacedResourceTimeline(ctx, ns, containerFields.PodName) - bindingPath := k8saudit.MustK8sSubresourceTimeline(ctx, podPath, "binding") + podPath, bindingPath := mustPodAndBindingTimelinePaths(ctx, clusterName, containerFields.Namespace, containerFields.PodName) nodeNameChanged := state == nil || state.LastNodeName != nodeFields.NodeName labelsChanged := state == nil || !maps.Equal(state.LastLabels, nodeFields.PodLabels) @@ -247,46 +240,33 @@ func (m *containerLogPodPhaseTimelineMapper) ProcessLogByGroup(ctx context.Conte podPhasePath = mustPodPhaseTimelinePath(ctx, clusterName, nodeFields.NodeName, containerFields.Namespace, containerFields.PodName, uid) } - // Check if audit log has already written to the Pod, its binding, or its phase timeline. - // When CAI knows about the Pod, a revision on podPath may originate from CAI rather than audit logs. + // Check if audit log has already recorded creation times for the Pod, its binding, or its phase timeline. if state == nil { - builder := khictx.MustGetValue(ctx, inspectioncore.Builder) - hasBindingRevision := builder.TimelineAccumulator.HasRevision(bindingPath) - hasPodPhaseRevision := builder.TimelineAccumulator.HasRevision(podPhasePath) - hasAuditPodRevision := !hasInitialState && builder.TimelineAccumulator.HasRevision(podPath) + timelineCreationTimes := coretask.GetTaskResult(ctx, k8saudit.TimelineCreationTimeInventoryTaskID.Ref()) + hasPodCreationTime := len(timelineCreationTimes[podPath]) > 0 + hasBindingCreationTime := len(timelineCreationTimes[bindingPath]) > 0 + hasPodPhaseCreationTime := len(timelineCreationTimes[podPhasePath]) > 0 - if hasBindingRevision || hasPodPhaseRevision || hasAuditPodRevision { + if hasPodCreationTime || hasBindingCreationTime || hasPodPhaseCreationTime { return nil, &containerLogPodPhaseMapperState{AuditLogFound: true}, nil } } - labels := map[string]any{} - for k, v := range nodeFields.PodLabels { - labels[k] = v - } - - podManifest := map[string]any{ - "apiVersion": "v1", - "kind": "Pod", - "metadata": map[string]any{ - "name": containerFields.PodName, - "namespace": containerFields.Namespace, - "labels": labels, - }, - "spec": map[string]any{ - "nodeName": nodeFields.NodeName, - }, - } - podNode, err := structured.FromGoValue(podManifest, &structured.AlphabeticalGoMapKeyOrderProvider{}) + podNode, err := makePodManifestNode(containerFields.Namespace, containerFields.PodName, nodeFields.NodeName, nodeFields.PodLabels) if err != nil { return nil, state, fmt.Errorf("failed to generate pod manifest: %w", err) } + changedTime := time.Unix(0, 0) + if state != nil { + changedTime = l.Timestamp + } + cs := khifilev6.NewTimelineChangeSet(l) if nodeNameChanged { cs.AddRevision(podPhasePath, &khifilev6.StagingRevision{ - ChangedTime: time.Unix(0, 0), + ChangedTime: changedTime, ResourceBody: podNode, Principal: "N/A", VerbType: k8saudit.VerbUnknown, @@ -297,7 +277,7 @@ func (m *containerLogPodPhaseTimelineMapper) ProcessLogByGroup(ctx context.Conte // Only supplement Pod and Binding revisions if CAI does not know about this Pod. if !hasInitialState { cs.AddRevision(podPath, &khifilev6.StagingRevision{ - ChangedTime: time.Unix(0, 0), + ChangedTime: changedTime, ResourceBody: podNode, Principal: "N/A", VerbType: k8saudit.VerbUnknown, @@ -305,24 +285,12 @@ func (m *containerLogPodPhaseTimelineMapper) ProcessLogByGroup(ctx context.Conte }) if nodeNameChanged { - bindingManifest := map[string]any{ - "apiVersion": "v1", - "kind": "Binding", - "metadata": map[string]any{ - "name": containerFields.PodName, - "namespace": containerFields.Namespace, - }, - "target": map[string]any{ - "kind": "Node", - "name": nodeFields.NodeName, - }, - } - bindingNode, err := structured.FromGoValue(bindingManifest, &structured.AlphabeticalGoMapKeyOrderProvider{}) + bindingNode, err := makeBindingManifestNode(containerFields.Namespace, containerFields.PodName, nodeFields.NodeName) if err != nil { return nil, state, fmt.Errorf("failed to generate binding manifest: %w", err) } cs.AddRevision(bindingPath, &khifilev6.StagingRevision{ - ChangedTime: time.Unix(0, 0), + ChangedTime: changedTime, ResourceBody: bindingNode, Principal: "N/A", VerbType: k8saudit.VerbUnknown, @@ -334,6 +302,56 @@ func (m *containerLogPodPhaseTimelineMapper) ProcessLogByGroup(ctx context.Conte return cs, nextState, nil } +// mustPodAndBindingTimelinePaths resolves the timeline paths for a Pod and its binding subresource. +func mustPodAndBindingTimelinePaths(ctx context.Context, clusterName, namespace, podName string) (*khifilev6.TimelinePath, *khifilev6.TimelinePath) { + cluster := k8saudit.MustK8sClusterTimeline(ctx, clusterName) + api := k8saudit.MustK8sAPIVersionTimeline(ctx, cluster, "core/v1") + kind := k8saudit.MustK8sKindTimeline(ctx, api, "pod") + ns := k8saudit.MustK8sNamespaceTimeline(ctx, kind, namespace) + podPath := k8saudit.MustK8sNamespacedResourceTimeline(ctx, ns, podName) + bindingPath := k8saudit.MustK8sSubresourceTimeline(ctx, podPath, "binding") + return podPath, bindingPath +} + +// makePodManifestNode constructs a structured YAML node for a synthesized Pod manifest. +func makePodManifestNode(namespace, podName, nodeName string, podLabels map[string]string) (structured.Node, error) { + labels := map[string]any{} + for k, v := range podLabels { + labels[k] = v + } + + podManifest := map[string]any{ + "apiVersion": "v1", + "kind": "Pod", + "metadata": map[string]any{ + "name": podName, + "namespace": namespace, + "labels": labels, + }, + "spec": map[string]any{ + "nodeName": nodeName, + }, + } + return structured.FromGoValue(podManifest, &structured.AlphabeticalGoMapKeyOrderProvider{}) +} + +// makeBindingManifestNode constructs a structured YAML node for a synthesized Binding manifest. +func makeBindingManifestNode(namespace, podName, nodeName string) (structured.Node, error) { + bindingManifest := map[string]any{ + "apiVersion": "v1", + "kind": "Binding", + "metadata": map[string]any{ + "name": podName, + "namespace": namespace, + }, + "target": map[string]any{ + "kind": "Node", + "name": nodeName, + }, + } + return structured.FromGoValue(bindingManifest, &structured.AlphabeticalGoMapKeyOrderProvider{}) +} + func mustPodPhaseTimelinePath(ctx context.Context, clusterName, nodeName, namespace, podName, uid string) *khifilev6.TimelinePath { cluster := k8saudit.MustK8sClusterTimeline(ctx, clusterName) api := k8saudit.MustK8sAPIVersionTimeline(ctx, cluster, "core/v1") diff --git a/pkg/task/inspection/googlecloud/k8scontainer/impl/mapper_test.go b/pkg/task/inspection/googlecloud/k8scontainer/impl/mapper_test.go index fcb22de72..8084ce81d 100644 --- a/pkg/task/inspection/googlecloud/k8scontainer/impl/mapper_test.go +++ b/pkg/task/inspection/googlecloud/k8scontainer/impl/mapper_test.go @@ -353,13 +353,12 @@ func TestPodPhaseTimelineMapper_ProcessLogByGroup(t *testing.T) { } testCases := []struct { - name string - inputLogs []*log.Log - cluster k8scommon.GoogleCloudClusterIdentity - initialStateProvider *mockInitialResourceStateProvider - flushToAccumulator bool - setup func() - assert func(t *testing.T, ctx context.Context, css []*khifilev6.TimelineChangeSet) + name string + inputLogs []*log.Log + cluster k8scommon.GoogleCloudClusterIdentity + initialStateProvider *mockInitialResourceStateProvider + timelineCreationTimes func() k8saudit.TimelineCreationTimes + assert func(t *testing.T, ctx context.Context, css []*khifilev6.TimelineChangeSet) }{ { name: "skipped because NodeName is empty", @@ -446,9 +445,9 @@ func TestPodPhaseTimelineMapper_ProcessLogByGroup(t *testing.T) { cluster: k8scommon.GoogleCloudClusterIdentity{ ClusterName: "test-cluster", }, - setup: func() { + timelineCreationTimes: func() k8saudit.TimelineCreationTimes { auditPodPath := k8saudit.MustK8sNamespacedResourceTimeline(ctx, namespaceTimeline, "test-pod-audit") - builder.TimelineAccumulator.AddTestRevision(auditPodPath) + return k8saudit.TimelineCreationTimes{auditPodPath: []time.Time{time.Date(2026, 5, 26, 11, 0, 0, 0, time.UTC)}} }, assert: func(t *testing.T, ctx context.Context, css []*khifilev6.TimelineChangeSet) { if css[0] != nil { @@ -474,10 +473,10 @@ func TestPodPhaseTimelineMapper_ProcessLogByGroup(t *testing.T) { cluster: k8scommon.GoogleCloudClusterIdentity{ ClusterName: "test-cluster", }, - setup: func() { + timelineCreationTimes: func() k8saudit.TimelineCreationTimes { bindingPodPath := k8saudit.MustK8sNamespacedResourceTimeline(ctx, namespaceTimeline, "test-pod-binding") subresourcePath := k8saudit.MustK8sSubresourceTimeline(ctx, bindingPodPath, "binding") - builder.TimelineAccumulator.AddTestRevision(subresourcePath) + return k8saudit.TimelineCreationTimes{subresourcePath: []time.Time{time.Date(2026, 5, 26, 11, 0, 0, 0, time.UTC)}} }, assert: func(t *testing.T, ctx context.Context, css []*khifilev6.TimelineChangeSet) { if css[0] != nil { @@ -611,7 +610,7 @@ func TestPodPhaseTimelineMapper_ProcessLogByGroup(t *testing.T) { // Second log generates only Pod (labels changed, node remained same) testchangeset.AssertTimeline(t, css[1]). HasRevision(podPath, &khifilev6.StagingRevision{ - ChangedTime: time.Unix(0, 0), + ChangedTime: time.Date(2026, 5, 26, 12, 0, 1, 0, time.UTC), ResourceBody: makePodNode("test-node", map[string]string{"a": "2"}), Principal: "N/A", VerbType: k8saudit.VerbUnknown, @@ -691,21 +690,21 @@ func TestPodPhaseTimelineMapper_ProcessLogByGroup(t *testing.T) { // Second log generates all 3 on node 2 testchangeset.AssertTimeline(t, css[1]). HasRevision(expectedPath2, &khifilev6.StagingRevision{ - ChangedTime: time.Unix(0, 0), + ChangedTime: time.Date(2026, 5, 26, 12, 0, 1, 0, time.UTC), ResourceBody: makePodNode("test-node-2", nil), Principal: "N/A", VerbType: k8saudit.VerbUnknown, StateType: k8saudit.RevisionStatePodPhaseUnknown, }, nodeComparer). HasRevision(podPath, &khifilev6.StagingRevision{ - ChangedTime: time.Unix(0, 0), + ChangedTime: time.Date(2026, 5, 26, 12, 0, 1, 0, time.UTC), ResourceBody: makePodNode("test-node-2", nil), Principal: "N/A", VerbType: k8saudit.VerbUnknown, StateType: k8saudit.RevisionStateK8sResourceExistingLogNotFound, }, nodeComparer). HasRevision(bindingPath, &khifilev6.StagingRevision{ - ChangedTime: time.Unix(0, 0), + ChangedTime: time.Date(2026, 5, 26, 12, 0, 1, 0, time.UTC), ResourceBody: makeBindingNode("test-node-2"), Principal: "N/A", VerbType: k8saudit.VerbUnknown, @@ -786,10 +785,6 @@ func TestPodPhaseTimelineMapper_ProcessLogByGroup(t *testing.T) { }: {}, }, }, - setup: func() { - caiPodPath := k8saudit.MustK8sNamespacedResourceTimeline(ctx, namespaceTimeline, "test-pod-cai-rev") - builder.TimelineAccumulator.AddTestRevision(caiPodPath) - }, assert: func(t *testing.T, ctx context.Context, css []*khifilev6.TimelineChangeSet) { caiPodPhasePath := mustPodPhaseTimelinePath(ctx, "test-cluster", "test-node", "test-namespace", "test-pod-cai-rev", "unknown") caiPodPath := k8saudit.MustK8sNamespacedResourceTimeline(ctx, namespaceTimeline, "test-pod-cai-rev") @@ -807,6 +802,44 @@ func TestPodPhaseTimelineMapper_ProcessLogByGroup(t *testing.T) { HasNoRevision(caiBindingPath) }, }, + { + name: "CAI known Pod with audit log Pod revision: skipped because audit log takes precedence", + inputLogs: []*log.Log{ + testlog.NewMockLog( + k8scontainer.K8sContainerLogFieldSet{ + Namespace: "test-namespace", + PodName: "test-pod-cai-audit", + ContainerName: "test-container", + }, + k8scontainer.GCPContainerLogNodeNameLabelFieldSet{ + NodeName: "test-node", + }, + time.Date(2026, 5, 26, 12, 0, 0, 0, time.UTC), + ), + }, + cluster: k8scommon.GoogleCloudClusterIdentity{ + ClusterName: "test-cluster", + }, + initialStateProvider: &mockInitialResourceStateProvider{ + states: map[k8saudit.ResourceIdentity]*structured.NodeReader{ + { + APIVersion: "core/v1", + Kind: "pod", + Namespace: "test-namespace", + Name: "test-pod-cai-audit", + }: {}, + }, + }, + timelineCreationTimes: func() k8saudit.TimelineCreationTimes { + caiPodPath := k8saudit.MustK8sNamespacedResourceTimeline(ctx, namespaceTimeline, "test-pod-cai-audit") + return k8saudit.TimelineCreationTimes{caiPodPath: []time.Time{time.Date(2026, 5, 26, 11, 0, 0, 0, time.UTC)}} + }, + assert: func(t *testing.T, ctx context.Context, css []*khifilev6.TimelineChangeSet) { + if css[0] != nil { + t.Errorf("expected cs to be nil, got %v", css[0]) + } + }, + }, { name: "CAI known Pod: subsequent log with only label changes returns nil changeset", inputLogs: []*log.Log{ @@ -888,10 +921,10 @@ func TestPodPhaseTimelineMapper_ProcessLogByGroup(t *testing.T) { }: {}, }, }, - setup: func() { + timelineCreationTimes: func() k8saudit.TimelineCreationTimes { caiPodPath := k8saudit.MustK8sNamespacedResourceTimeline(ctx, namespaceTimeline, "test-pod-cai-binding") subresourcePath := k8saudit.MustK8sSubresourceTimeline(ctx, caiPodPath, "binding") - builder.TimelineAccumulator.AddTestRevision(subresourcePath) + return k8saudit.TimelineCreationTimes{subresourcePath: []time.Time{time.Date(2026, 5, 26, 11, 0, 0, 0, time.UTC)}} }, assert: func(t *testing.T, ctx context.Context, css []*khifilev6.TimelineChangeSet) { if css[0] != nil { @@ -927,9 +960,9 @@ func TestPodPhaseTimelineMapper_ProcessLogByGroup(t *testing.T) { }: makeCAIPodNodeReader("cai-pod-uid-456"), }, }, - setup: func() { + timelineCreationTimes: func() k8saudit.TimelineCreationTimes { existingPhasePath := mustPodPhaseTimelinePath(ctx, "test-cluster", "test-node", "test-namespace", "test-pod-cai-phase", "cai-pod-uid-456") - builder.TimelineAccumulator.AddTestRevision(existingPhasePath) + return k8saudit.TimelineCreationTimes{existingPhasePath: []time.Time{time.Date(2026, 5, 26, 11, 0, 0, 0, time.UTC)}} }, assert: func(t *testing.T, ctx context.Context, css []*khifilev6.TimelineChangeSet) { if css[0] != nil { @@ -937,76 +970,11 @@ func TestPodPhaseTimelineMapper_ProcessLogByGroup(t *testing.T) { } }, }, - { - name: "non-CAI Pod with flushed accumulator between logs: subsequent log still maps node change", - inputLogs: []*log.Log{ - testlog.NewMockLog( - k8scontainer.K8sContainerLogFieldSet{ - Namespace: "test-namespace", - PodName: "test-pod-flushed", - ContainerName: "test-container", - }, - k8scontainer.GCPContainerLogNodeNameLabelFieldSet{ - NodeName: "test-node-1", - }, - time.Date(2026, 5, 26, 12, 0, 0, 0, time.UTC), - ), - testlog.NewMockLog( - k8scontainer.K8sContainerLogFieldSet{ - Namespace: "test-namespace", - PodName: "test-pod-flushed", - ContainerName: "test-container", - }, - k8scontainer.GCPContainerLogNodeNameLabelFieldSet{ - NodeName: "test-node-2", - }, - time.Date(2026, 5, 26, 12, 0, 1, 0, time.UTC), - ), - }, - cluster: k8scommon.GoogleCloudClusterIdentity{ - ClusterName: "test-cluster", - }, - flushToAccumulator: true, - assert: func(t *testing.T, ctx context.Context, css []*khifilev6.TimelineChangeSet) { - if len(css) != 2 { - t.Fatalf("expected 2 changesets, got %d", len(css)) - } - podPhasePath2 := mustPodPhaseTimelinePath(ctx, "test-cluster", "test-node-2", "test-namespace", "test-pod-flushed", "unknown") - flushedPodPath := k8saudit.MustK8sNamespacedResourceTimeline(ctx, namespaceTimeline, "test-pod-flushed") - flushedBindingPath := k8saudit.MustK8sSubresourceTimeline(ctx, flushedPodPath, "binding") - - testchangeset.AssertTimeline(t, css[1]). - HasRevision(podPhasePath2, &khifilev6.StagingRevision{ - ChangedTime: time.Unix(0, 0), - ResourceBody: makeNamedPodNode("test-pod-flushed", "test-node-2", nil), - Principal: "N/A", - VerbType: k8saudit.VerbUnknown, - StateType: k8saudit.RevisionStatePodPhaseUnknown, - }, nodeComparer). - HasRevision(flushedPodPath, &khifilev6.StagingRevision{ - ChangedTime: time.Unix(0, 0), - ResourceBody: makeNamedPodNode("test-pod-flushed", "test-node-2", nil), - Principal: "N/A", - VerbType: k8saudit.VerbUnknown, - StateType: k8saudit.RevisionStateK8sResourceExistingLogNotFound, - }, nodeComparer). - HasRevision(flushedBindingPath, &khifilev6.StagingRevision{ - ChangedTime: time.Unix(0, 0), - ResourceBody: makeNamedBindingNode("test-pod-flushed", "test-node-2"), - Principal: "N/A", - VerbType: k8saudit.VerbUnknown, - StateType: k8saudit.RevisionStateK8sResourceExistingLogNotFound, - }, nodeComparer) - }, - }, } mapper := &containerLogPodPhaseTimelineMapper{} for _, tc := range testCases { t.Run(tc.name, func(t *testing.T) { - if tc.setup != nil { - tc.setup() - } ctx := khictx.WithValue(t.Context(), inspectioncore.Builder, builder) ctx = tasktest.WithTaskResult(ctx, k8scontainer.ClusterIdentityTaskID.Ref(), tc.cluster) var provider k8saudit.InitialResourceStateProvider = &mockInitialResourceStateProvider{} @@ -1015,6 +983,12 @@ func TestPodPhaseTimelineMapper_ProcessLogByGroup(t *testing.T) { } ctx = tasktest.WithTaskResult(ctx, k8saudit.InitialResourceStateProviderRef, provider) + creationTimes := k8saudit.TimelineCreationTimes{} + if tc.timelineCreationTimes != nil { + creationTimes = tc.timelineCreationTimes() + } + ctx = tasktest.WithTaskResult(ctx, k8saudit.TimelineCreationTimeInventoryTaskID.Ref(), creationTimes) + var css []*khifilev6.TimelineChangeSet var state *containerLogPodPhaseMapperState for _, l := range tc.inputLogs { @@ -1022,11 +996,6 @@ func TestPodPhaseTimelineMapper_ProcessLogByGroup(t *testing.T) { if err != nil { t.Fatalf("ProcessLogByGroup() returned unexpected error: %v", err) } - if tc.flushToAccumulator && cs != nil { - cs.ForEachRevision(func(path *khifilev6.TimelinePath, _ []*khifilev6.StagingRevision) { - builder.TimelineAccumulator.AddTestRevision(path) - }) - } css = append(css, cs) state = nextState }