diff --git a/cmd/krply/record.go b/cmd/krply/record.go index 2b2ca1c..a5b09c8 100644 --- a/cmd/krply/record.go +++ b/cmd/krply/record.go @@ -2,6 +2,7 @@ package main import ( "context" + "errors" "fmt" "os" "os/signal" @@ -47,6 +48,9 @@ func init() { } func runRecord(cmd *cobra.Command, args []string) error { + if serverURL != "" { + return errors.New("record runs against a live cluster with a local store; --server is not supported (omit --server)") + } store, err := openStore(storePath) if err != nil { return err diff --git a/internal/api/server.go b/internal/api/server.go index d92c9b8..87cdd1f 100644 --- a/internal/api/server.go +++ b/internal/api/server.go @@ -320,6 +320,9 @@ func (s *Server) handleObjectHistory(w http.ResponseWriter, r *http.Request) { s.internalError(w, "gaps", err) return } + // Scope gap markers to the requested window so an out-of-window historical + // gap does not make a gap-free window look incomplete. + gaps = filterByTime(gaps, since, until) history := queryv1.ObjectHistory{ ClusterID: ref.ClusterID, diff --git a/internal/event/stream_test.go b/internal/event/stream_test.go new file mode 100644 index 0000000..922e5dc --- /dev/null +++ b/internal/event/stream_test.go @@ -0,0 +1,66 @@ +package event + +import "testing" + +func TestStreamIDSelectorRoundTrip(t *testing.T) { + // Label selectors legitimately contain '/' (prefix/name label keys). The + // ID must round-trip through StreamID without loss or ambiguity. + stream := Stream{ + ClusterID: "cluster-abc", + Group: "apps", + Version: "v1", + Resource: "deployments", + Namespace: "default", + Selector: "app.kubernetes.io/name=nginx,env=prod", + } + id := stream.ID() + parsed, err := StreamID(id) + if err != nil { + t.Fatalf("StreamID(%q): %v", id, err) + } + if parsed != stream { + t.Fatalf("StreamID round-trip = %+v, want %+v", parsed, stream) + } +} + +func TestStreamIDSelectorAmbiguity(t *testing.T) { + // {Namespace: e, Selector: f/g} must not collide with {Namespace: e/f, + // Selector: g}: the selector's slash is escaped, so the two IDs differ. + a := Stream{ClusterID: "c", Group: "g", Version: "v", Resource: "r", Namespace: "e", Selector: "f/g"} + b := Stream{ClusterID: "c", Group: "g", Version: "v", Resource: "r", Namespace: "e/f", Selector: "g"} + if a.ID() == b.ID() { + t.Fatalf("streams with slash in selector collided: %q", a.ID()) + } + if _, err := StreamID(a.ID()); err != nil { + t.Fatalf("StreamID(%q): %v", a.ID(), err) + } +} + +func TestStreamIDNoSelector(t *testing.T) { + stream := Stream{ClusterID: "c", Group: "", Version: "v1", Resource: "pods", Namespace: "default"} + id := stream.ID() + parsed, err := StreamID(id) + if err != nil { + t.Fatalf("StreamID(%q): %v", id, err) + } + if parsed != stream { + t.Fatalf("round-trip = %+v, want %+v", parsed, stream) + } +} + +func TestEventIDLiveVsSynthetic(t *testing.T) { + stream := Stream{ClusterID: "c", Group: "", Version: "v1", Resource: "pods", Namespace: "default"} + ref := ResourceRef{Namespace: "default", Name: "p", UID: "u", ResourceVersion: "100"} + // Live events (observedAt 0) dedup to the same key. + if EventID(stream, ref, WatchAdded, 0) != EventID(stream, ref, WatchAdded, 0) { + t.Fatal("live event_id not deterministic") + } + // Synthetic relist events (non-zero observedAt) differ from live events and + // from each other, so an unchanged re-listed object survives dedup. + if EventID(stream, ref, WatchAdded, 0) == EventID(stream, ref, WatchAdded, 1) { + t.Fatal("synthetic event_id must differ from the live key") + } + if EventID(stream, ref, WatchAdded, 1) == EventID(stream, ref, WatchAdded, 2) { + t.Fatal("synthetic event_ids must differ across relists") + } +} diff --git a/internal/materialize/materialize_test.go b/internal/materialize/materialize_test.go index 45e6753..f6fb01d 100644 --- a/internal/materialize/materialize_test.go +++ b/internal/materialize/materialize_test.go @@ -312,3 +312,185 @@ func TestDiff(t *testing.T) { t.Fatalf("C changes = %+v, want single whole-object Added", cDiff.Changes) } } + +// TestGapHealingAndBaselineReset verifies that a BASELINE after a gap heals +// the open gap and resets object state, so an object deleted during the gap +// disappears while a surviving object is re-materialized from the synthetic +// ADDED that follows the baseline. +func TestGapHealingAndBaselineReset(t *testing.T) { + store := storage.NewInMemory() + mat := NewMaterializer(store) + ctx := context.Background() + + cluster := "c1" + sid := event.Stream{ClusterID: cluster, Group: "apps", Version: "v1", Resource: "deployments", Namespace: "default"}.ID() + base := time.Date(2026, 8, 1, 0, 0, 0, 0, time.UTC) + + appendRec(t, store, baseline(cluster, sid, base)) + appendRec(t, store, ev(cluster, sid, event.WatchAdded, deployRef("A"), deployObject("A", "nginx:1.20", "1"), base.Add(time.Minute))) + appendRec(t, store, gap(cluster, sid, base.Add(2*time.Minute))) + // After the gap: B survived, A was deleted during the gap. + appendRec(t, store, baseline(cluster, sid, base.Add(3*time.Minute))) + appendRec(t, store, ev(cluster, sid, event.WatchAdded, deployRef("B"), deployObject("B", "nginx:1.30", "2"), base.Add(3*time.Minute+time.Second))) + + res, err := mat.StreamState(ctx, sid, base.Add(4*time.Minute)) + if err != nil { + t.Fatalf("StreamState: %v", err) + } + if res.HasGaps { + t.Fatalf("HasGaps = true, want false (baseline healed the gap)") + } + names := map[string]bool{} + for _, o := range res.Objects { + names[o.Name] = true + } + if !names["B"] { + t.Fatalf("object B missing from state after healed gap: %+v", names) + } + if names["A"] { + t.Fatalf("object A (deleted during the gap) still present: %+v", names) + } +} + +// TestDiffIgnoreIsPathScoped verifies that server-metadata fields are ignored +// only under metadata, while user fields with the same names are kept. +func TestDiffIgnoreIsPathScoped(t *testing.T) { + store := storage.NewInMemory() + mat := NewMaterializer(store) + ctx := context.Background() + + cluster := "c1" + sid := event.Stream{ClusterID: cluster, Group: "", Version: "v1", Resource: "configmaps", Namespace: "default"}.ID() + base := time.Date(2026, 8, 1, 0, 0, 0, 0, time.UTC) + + mk := func(rv string, dataStatus string, dataUID string) json.RawMessage { + obj := map[string]any{ + "apiVersion": "v1", + "kind": "ConfigMap", + "metadata": map[string]any{ + "name": "app", "namespace": "default", + "uid": "u", "resourceVersion": rv, "generation": 1, + }, + "data": map[string]any{"status": dataStatus, "uid": dataUID}, + } + raw, _ := json.Marshal(obj) + return raw + } + + ref := event.ResourceRef{Namespace: "default", Name: "app", Kind: "ConfigMap"} + appendRec(t, store, baseline(cluster, sid, base)) + appendRec(t, store, ev(cluster, sid, event.WatchAdded, ref, mk("1", "up", "v1"), base.Add(time.Minute))) + appendRec(t, store, ev(cluster, sid, event.WatchModified, ref, mk("2", "down", "v2"), base.Add(2*time.Minute))) + + res, err := mat.Diff(ctx, cluster, "default", base.Add(90*time.Second), base.Add(150*time.Second)) + if err != nil { + t.Fatalf("Diff: %v", err) + } + if len(res.Changes) != 1 { + t.Fatalf("changed objects = %d, want 1", len(res.Changes)) + } + paths := map[string]bool{} + for _, ch := range res.Changes[0].Changes { + paths[ch.Path] = true + } + if !paths["data.status"] { + t.Fatalf("data.status change missing (user field dropped): %v", paths) + } + if !paths["data.uid"] { + t.Fatalf("data.uid change missing (user field dropped): %v", paths) + } + for _, p := range []string{"metadata.uid", "metadata.resourceVersion", "metadata.generation", "status"} { + if paths[p] { + t.Fatalf("server metadata field %q leaked into the diff: %v", p, paths) + } + } +} + +// TestDiffNestedAddedRemovedFlags verifies nested adds/removes carry the +// Added/Removed flags so the UI and CLI can render them correctly. +func TestDiffNestedAddedRemovedFlags(t *testing.T) { + store := storage.NewInMemory() + mat := NewMaterializer(store) + ctx := context.Background() + + cluster := "c1" + sid := event.Stream{ClusterID: cluster, Group: "apps", Version: "v1", Resource: "deployments", Namespace: "default"}.ID() + base := time.Date(2026, 8, 1, 0, 0, 0, 0, time.UTC) + ref := deployRef("A") + + mk := func(rv string, label string) json.RawMessage { + obj := map[string]any{ + "apiVersion": "apps/v1", + "kind": "Deployment", + "metadata": map[string]any{"name": "A", "namespace": "default", "labels": map[string]any{"team": label}}, + } + raw, _ := json.Marshal(obj) + return raw + } + + appendRec(t, store, baseline(cluster, sid, base)) + appendRec(t, store, ev(cluster, sid, event.WatchAdded, ref, mk("1", "a"), base.Add(time.Minute))) + appendRec(t, store, ev(cluster, sid, event.WatchModified, ref, mk("2", "b"), base.Add(2*time.Minute))) + + res, err := mat.Diff(ctx, cluster, "default", base.Add(90*time.Second), base.Add(150*time.Second)) + if err != nil { + t.Fatalf("Diff: %v", err) + } + if len(res.Changes) != 1 || len(res.Changes[0].Changes) != 1 { + t.Fatalf("expected 1 change, got %+v", res.Changes) + } + ch := res.Changes[0].Changes[0] + if ch.Path != "metadata.labels.team" || ch.Added || ch.Removed || ch.Before != "a" || ch.After != "b" { + t.Fatalf("changed field = %+v, want metadata.labels.team a -> b", ch) + } +} + +// TestSnapshotWritesRecordAndWatermarks verifies Snapshot persists a +// TypeSnapshot journal record and per-stream watermarks. +func TestSnapshotWritesRecordAndWatermarks(t *testing.T) { + store := storage.NewInMemory() + mat := NewMaterializer(store) + ctx := context.Background() + + cluster := "c1" + sid := event.Stream{ClusterID: cluster, Group: "", Version: "v1", Resource: "configmaps", Namespace: "default"}.ID() + base := time.Date(2026, 8, 1, 0, 0, 0, 0, time.UTC) + + appendRec(t, store, baseline(cluster, sid, base)) + appendRec(t, store, ev(cluster, sid, event.WatchAdded, event.ResourceRef{Namespace: "default", Name: "cm", Kind: "ConfigMap", ResourceVersion: "42"}, configMapObject("cm", "42"), base.Add(time.Minute))) + + snap, err := mat.Snapshot(ctx, cluster, base.Add(2*time.Minute), "test") + if err != nil { + t.Fatalf("Snapshot: %v", err) + } + if !snap.Complete { + t.Fatalf("Snapshot not complete: %s", snap.Warning) + } + if len(snap.Watermarks) != 1 { + t.Fatalf("watermarks = %d, want 1", len(snap.Watermarks)) + } + if snap.Watermarks[0].LastResourceVersion != "42" || snap.Watermarks[0].LastObservedAt.IsZero() { + t.Fatalf("watermark = %+v", snap.Watermarks[0]) + } + + found := false + for _, rec := range storeEvents(ctx, store) { + if rec.Type == event.TypeSnapshot { + found = true + if rec.Snapshot == nil || rec.Snapshot.Name != "test" { + t.Fatalf("snapshot record missing name: %+v", rec.Snapshot) + } + } + } + if !found { + t.Fatal("no TypeSnapshot journal record was written") + } +} + +func storeEvents(ctx context.Context, store storage.Store) []event.Record { + recs, err := store.Events(ctx, storage.EventFilter{}) + if err != nil { + return nil + } + return recs +} diff --git a/internal/replay/replay_test.go b/internal/replay/replay_test.go index 481191d..630eaf2 100644 --- a/internal/replay/replay_test.go +++ b/internal/replay/replay_test.go @@ -318,3 +318,96 @@ func TestPlanNamespaceMapping(t *testing.T) { t.Fatalf("plan warnings %v lack a namespace mapping note", plan.Warnings) } } + +func TestSanitizeServiceNodePortAndHeadless(t *testing.T) { + nodePort := map[string]any{ + "apiVersion": "v1", + "kind": "Service", + "metadata": map[string]any{"name": "np", "namespace": "default"}, + "spec": map[string]any{ + "type": "NodePort", + "clusterIP": "10.96.0.5", + "clusterIPs": []any{"10.96.0.5"}, + "ports": []any{map[string]any{"port": 80, "nodePort": 30080}}, + }, + } + clean, _ := sanitizeObject(nodePort, "Service", DefaultPolicy()) + spec := clean["spec"].(map[string]any) + if _, ok := spec["clusterIP"]; ok { + t.Fatalf("NodePort service kept clusterIP: %+v", spec) + } + if _, ok := spec["clusterIPs"]; ok { + t.Fatalf("NodePort service kept clusterIPs: %+v", spec) + } + ports := spec["ports"].([]any) + pm := ports[0].(map[string]any) + if _, ok := pm["nodePort"]; ok { + t.Fatalf("NodePort service kept nodePort: %+v", pm) + } + + headless := map[string]any{ + "apiVersion": "v1", + "kind": "Service", + "metadata": map[string]any{"name": "hl", "namespace": "default"}, + "spec": map[string]any{"type": "ClusterIP", "clusterIP": "None"}, + } + clean2, _ := sanitizeObject(headless, "Service", DefaultPolicy()) + if got := clean2["spec"].(map[string]any)["clusterIP"]; got != "None" { + t.Fatalf("headless service clusterIP = %v, want None", got) + } +} + +func TestApplyRequiresDryRun(t *testing.T) { + plan := &Plan{ + ID: "plan-test", + Status: "planned", + Objects: []PlanObject{{ + Namespace: "default", Name: "x", Kind: "ConfigMap", + Object: map[string]any{"apiVersion": "v1", "kind": "ConfigMap", "metadata": map[string]any{"name": "x"}}, + }}, + } + _, err := plan.Apply(context.Background(), "", "", true) + if err == nil { + t.Fatal("Apply without a successful dry run must be refused") + } +} + +func TestApplyItemsReportsSkipped(t *testing.T) { + plan := &Plan{ + Objects: []PlanObject{ + {Namespace: "default", Name: "known", Kind: "ConfigMap", Object: map[string]any{"apiVersion": "v1", "kind": "ConfigMap", "metadata": map[string]any{"name": "known"}}}, + {Namespace: "default", Name: "unknown", Kind: "MysteryKind", Object: map[string]any{"apiVersion": "acme.io/v1", "kind": "MysteryKind", "metadata": map[string]any{"name": "unknown"}}}, + }, + } + items, skipped := plan.applyItems() + if len(items) != 1 { + t.Fatalf("applyItems = %d items, want 1", len(items)) + } + if len(skipped) != 1 || skipped[0].Kind != "MysteryKind" { + t.Fatalf("skipped = %+v, want 1 MysteryKind", skipped) + } +} + +func TestGVRForCoversPolicyKinds(t *testing.T) { + for _, kind := range []string{"Secret", "Role", "ClusterRole", "RoleBinding", "ClusterRoleBinding", "Job", "CronJob", "Pod", "PersistentVolume", "PersistentVolumeClaim", "StorageClass", "MutatingWebhookConfiguration", "ValidatingWebhookConfiguration", "CustomResourceDefinition"} { + apiVersion := "v1" + switch kind { + case "Role", "RoleBinding", "ClusterRole", "ClusterRoleBinding": + apiVersion = "rbac.authorization.k8s.io/v1" + case "Job", "CronJob": + apiVersion = "batch/v1" + case "StorageClass": + apiVersion = "storage.k8s.io/v1" + case "PersistentVolume", "PersistentVolumeClaim": + apiVersion = "v1" + case "MutatingWebhookConfiguration", "ValidatingWebhookConfiguration": + apiVersion = "admissionregistration.k8s.io/v1" + case "CustomResourceDefinition": + apiVersion = "apiextensions.k8s.io/v1" + } + obj := map[string]any{"apiVersion": apiVersion, "kind": kind, "metadata": map[string]any{"name": "x"}} + if _, err := gvrFor(obj, kind); err != nil { + t.Fatalf("gvrFor(%s): %v", kind, err) + } + } +} diff --git a/internal/storage/time_test.go b/internal/storage/time_test.go new file mode 100644 index 0000000..18d6b50 --- /dev/null +++ b/internal/storage/time_test.go @@ -0,0 +1,121 @@ +package storage + +import ( + "context" + "path/filepath" + "testing" + "time" + + "github.com/krply/krply/internal/event" +) + +// TestSQLiteSubSecondBoundaries verifies that since/until and ObjectAt compare +// correctly across whole-second and sub-second boundaries. The old +// variable-width RFC3339Nano storage made lexicographic comparison wrong at +// sub-second precision; the fixed-width format must not. +func TestSQLiteSubSecondBoundaries(t *testing.T) { + store, err := NewSQLiteStore(filepath.Join(t.TempDir(), "t.db")) + if err != nil { + t.Fatalf("NewSQLiteStore: %v", err) + } + defer store.Close() + + ctx := context.Background() + sid := event.Stream{ClusterID: "c", Group: "apps", Version: "v1", Resource: "deployments", Namespace: "default"}.ID() + base := time.Date(2026, 8, 1, 0, 0, 0, 0, time.UTC) + + // An event exactly at a whole second and one half a second later. + whole := &event.Record{ + ClusterID: "c", StreamID: sid, Type: event.TypeEvent, EventID: "w", + ObservedAt: base, + WatchType: event.WatchAdded, + Resource: event.ResourceRef{Namespace: "default", Name: "a", Kind: "Deployment", ResourceVersion: "1"}, + Object: []byte(`{}`), + } + half := &event.Record{ + ClusterID: "c", StreamID: sid, Type: event.TypeEvent, EventID: "h", + ObservedAt: base.Add(500 * time.Millisecond), + WatchType: event.WatchAdded, + Resource: event.ResourceRef{Namespace: "default", Name: "b", Kind: "Deployment", ResourceVersion: "2"}, + Object: []byte(`{}`), + } + for _, r := range []*event.Record{whole, half} { + if _, err := store.Append(ctx, r); err != nil { + t.Fatalf("append: %v", err) + } + } + + // since = whole second must include both events (both are at or after it). + recs, err := store.Events(ctx, EventFilter{ClusterID: "c", Since: base}) + if err != nil { + t.Fatalf("Events since whole second: %v", err) + } + if len(recs) != 2 { + t.Fatalf("Events since whole second = %d, want 2", len(recs)) + } + + // until = whole second + 200ms must include only the whole-second event. + recs, err = store.Events(ctx, EventFilter{ClusterID: "c", Until: base.Add(200 * time.Millisecond)}) + if err != nil { + t.Fatalf("Events until sub-second: %v", err) + } + if len(recs) != 1 || recs[0].EventID != "w" { + t.Fatalf("Events until sub-second = %+v, want only 'w'", recs) + } + + // ObjectAt at whole second + 200ms must return the whole-second event. + obj, err := store.ObjectAt(ctx, ObjectRef{ClusterID: "c", StreamID: sid, Namespace: "default", Name: "a"}, base.Add(200*time.Millisecond)) + if err != nil { + t.Fatalf("ObjectAt: %v", err) + } + if obj.Resource.ResourceVersion != "1" { + t.Fatalf("ObjectAt RV = %q, want 1", obj.Resource.ResourceVersion) + } +} + +// TestMemoryStoreDedupAndPagination verifies the in-memory store matches the +// SQLite store: event_id dedup, SinceSeq cursors, and Limit/Offset. +func TestMemoryStoreDedupAndPagination(t *testing.T) { + store := NewInMemory() + ctx := context.Background() + sid := event.Stream{ClusterID: "c", Group: "", Version: "v1", Resource: "pods", Namespace: "default"}.ID() + at := time.Date(2026, 8, 1, 0, 0, 0, 0, time.UTC) + + mk := func(id string) *event.Record { + return &event.Record{ + ClusterID: "c", StreamID: sid, Type: event.TypeEvent, EventID: id, + ObservedAt: at, WatchType: event.WatchAdded, + Resource: event.ResourceRef{Namespace: "default", Name: id, Kind: "Pod", ResourceVersion: id}, + Object: []byte(`{}`), + } + } + for _, r := range []*event.Record{mk("a"), mk("b"), mk("a"), mk("c"), mk("d")} { + if _, err := store.Append(ctx, r); err != nil { + t.Fatalf("append: %v", err) + } + } + // "a" was appended twice; dedup must collapse it to one record. + recs, err := store.Events(ctx, EventFilter{ClusterID: "c"}) + if err != nil { + t.Fatalf("Events: %v", err) + } + if len(recs) != 4 { + t.Fatalf("deduped events = %d, want 4 (a,b,c,d)", len(recs)) + } + + // Cursor pagination: page of 2 after cursor 0. + page1, err := store.Events(ctx, EventFilter{ClusterID: "c", Limit: 2}) + if err != nil { + t.Fatalf("page1: %v", err) + } + if len(page1) != 2 || page1[0].EventID != "a" || page1[1].EventID != "b" { + t.Fatalf("page1 = %+v", page1) + } + page2, err := store.Events(ctx, EventFilter{ClusterID: "c", Limit: 2, SinceSeq: page1[1].IngestSeq}) + if err != nil { + t.Fatalf("page2: %v", err) + } + if len(page2) != 2 || page2[0].EventID != "c" || page2[1].EventID != "d" { + t.Fatalf("page2 = %+v", page2) + } +} diff --git a/test/fixtures/demo.db b/test/fixtures/demo.db index 3bbd60e..d61b6aa 100644 Binary files a/test/fixtures/demo.db and b/test/fixtures/demo.db differ