diff --git a/test/e2e/e2e_test.go b/test/e2e/e2e_test.go new file mode 100644 index 0000000..cf05042 --- /dev/null +++ b/test/e2e/e2e_test.go @@ -0,0 +1,225 @@ +//go:build e2e + +// Package e2e_test runs the full record -> timeline -> snapshot -> replay-plan +// pipeline against the fake apiserver without a real cluster. +package e2e_test + +import ( + "context" + "testing" + "time" + + "k8s.io/apimachinery/pkg/runtime/schema" + "k8s.io/client-go/dynamic" + "k8s.io/client-go/kubernetes/scheme" + "k8s.io/client-go/rest" + + "github.com/krply/krply/internal/discovery" + "github.com/krply/krply/internal/event" + "github.com/krply/krply/internal/materialize" + "github.com/krply/krply/internal/replay" + "github.com/krply/krply/internal/storage" + "github.com/krply/krply/internal/watch" + "github.com/krply/krply/test/integration/fakeapiserver" +) + +const clusterID = "e2e" + +func e2eDynamicClient(t *testing.T, url string) dynamic.Interface { + t.Helper() + cfg := &rest.Config{ + Host: url, + ContentConfig: rest.ContentConfig{ + NegotiatedSerializer: scheme.Codecs.WithoutConversion(), + GroupVersion: &schema.GroupVersion{Group: "", Version: "v1"}, + }, + BearerToken: "ignored", + } + dyn, err := dynamic.NewForConfig(cfg) + if err != nil { + t.Fatalf("dynamic.NewForConfig: %v", err) + } + return dyn +} + +func e2eCollector(t *testing.T, dyn dynamic.Interface, store storage.Store) *watch.Collector { + t.Helper() + col, err := watch.NewCollector(watch.Config{ + ClusterID: clusterID, + Resources: []discovery.ResourceSpec{ + {APIGroup: "", Version: "v1", Resource: "configmaps", Kind: "ConfigMap", Namespace: "default"}, + }, + Store: store, + DynamicClient: dyn, + Bookmarks: true, + MinBackoff: 10 * time.Millisecond, + MaxBackoff: 100 * time.Millisecond, + }) + if err != nil { + t.Fatalf("watch.NewCollector: %v", err) + } + return col +} + +func waitFor(t *testing.T, timeout time.Duration, cond func() bool, msg string) { + t.Helper() + deadline := time.Now().Add(timeout) + for time.Now().Before(deadline) { + if cond() { + return + } + time.Sleep(10 * time.Millisecond) + } + t.Fatalf("timed out waiting for %s", msg) +} + +func eventsOf(ctx context.Context, store storage.Store) []event.Record { + recs, err := store.Events(ctx, storage.EventFilter{}) + if err != nil { + return nil + } + return recs +} + +func hasDelete(store storage.Store, name string) func() bool { + return func() bool { + for _, r := range eventsOf(context.Background(), store) { + if r.Type == event.TypeEvent && r.WatchType == event.WatchDeleted && r.Resource.Name == name { + return true + } + } + return false + } +} + +func TestEndToEndPipeline(t *testing.T) { + fake, err := fakeapiserver.NewFake(map[string]string{ + "cm-a": `{"metadata":{"name":"cm-a"},"data":{"key":"v1"}}`, + "cm-b": `{"metadata":{"name":"cm-b"},"data":{"key":"v1"}}`, + }) + if err != nil { + t.Fatalf("NewFake: %v", err) + } + defer fake.Close() + + store := storage.NewInMemory() + col := e2eCollector(t, e2eDynamicClient(t, fake.URL), store) + + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + done := make(chan error, 1) + go func() { done <- col.Run(ctx) }() + + waitFor(t, 5*time.Second, func() bool { + var baseline, sawA, sawB bool + for _, r := range eventsOf(ctx, store) { + switch { + case r.Type == event.TypeBaseline: + baseline = true + case r.Type == event.TypeEvent && r.Synthetic && r.WatchType == event.WatchAdded && r.Resource.Name == "cm-a": + sawA = true + case r.Type == event.TypeEvent && r.Synthetic && r.WatchType == event.WatchAdded && r.Resource.Name == "cm-b": + sawB = true + } + } + return baseline && sawA && sawB + }, "baseline + synthetic ADDED for cm-a and cm-b") + + fake.AddOrUpdate("cm-a", map[string]any{ + "metadata": map[string]any{"name": "cm-a"}, + "data": map[string]any{"key": "v2"}, + }, "") + fake.Delete("cm-b") + + waitFor(t, 5*time.Second, hasDelete(store, "cm-b"), "DELETED event for cm-b") + time.Sleep(50 * time.Millisecond) + + streams, err := store.Streams(ctx) + if err != nil { + t.Fatalf("Streams: %v", err) + } + if len(streams) != 1 { + t.Fatalf("streams = %d, want 1", len(streams)) + } + sm := streams[0] + if !sm.Available || sm.HasGaps || sm.GapCount != 0 { + t.Fatalf("stream meta = %+v, want available with no gaps", sm) + } + if sm.LastResourceVersion != "102" { + t.Fatalf("last resource version = %q, want 102", sm.LastResourceVersion) + } + + history, err := store.ObjectHistory(ctx, storage.ObjectRef{ + ClusterID: clusterID, + StreamID: sm.StreamID, + Namespace: "default", + Name: "cm-a", + }) + if err != nil { + t.Fatalf("ObjectHistory: %v", err) + } + if len(history) != 2 { + t.Fatalf("cm-a history = %d records, want 2 (ADDED, MODIFIED)", len(history)) + } + if history[0].WatchType != event.WatchAdded || !history[0].Synthetic { + t.Fatalf("first history record = %+v, want synthetic ADDED", history[0]) + } + if history[1].WatchType != event.WatchModified { + t.Fatalf("second history record watch type = %q, want MODIFIED", history[1].WatchType) + } + + mat := materialize.NewMaterializer(store) + at := time.Now().UTC().Add(2 * time.Second) + snap, err := mat.Snapshot(ctx, clusterID, at, "e2e-snap") + if err != nil { + t.Fatalf("Snapshot: %v", err) + } + if !snap.Complete { + t.Fatalf("snapshot complete = false, missing = %v, warning = %q", snap.Missing, snap.Warning) + } + if len(snap.Objects) != 1 { + t.Fatalf("snapshot objects = %d, want 1 (deleted cm-b must not be present)", len(snap.Objects)) + } + if snap.Objects[0].Name != "cm-a" || snap.Objects[0].Kind != "ConfigMap" { + t.Fatalf("snapshot object = %+v, want cm-a ConfigMap", snap.Objects[0]) + } + + pol := replay.DefaultPolicy() + pol.AllowGaps = true + planner := replay.NewPlanner(store, mat, pol) + plan, err := planner.Plan(ctx, clusterID, snap.ID, "default", "") + if err != nil { + t.Fatalf("Plan: %v", err) + } + if plan.Status != "planned" || plan.SnapshotID != snap.ID { + t.Fatalf("plan = %+v, want status planned for snapshot %s", plan, snap.ID) + } + if !plan.CoverageComplete { + t.Fatalf("plan coverage complete = false, want true: the only stream has a baseline and no gaps") + } + if len(plan.Objects) != 1 { + t.Fatalf("plan objects = %d, want 1", len(plan.Objects)) + } + po := plan.Objects[0] + if po.Name != "cm-a" || po.Kind != "ConfigMap" || po.Namespace != "default" { + t.Fatalf("plan object = %+v, want cm-a ConfigMap in default", po) + } + meta, _ := po.Object["metadata"].(map[string]any) + for _, key := range []string{"uid", "resourceVersion", "namespace"} { + if _, ok := meta[key]; ok { + t.Fatalf("plan object metadata still contains %q: %+v", key, meta) + } + } + if _, ok := po.Object["status"]; ok { + t.Fatalf("plan object still contains status") + } + data, _ := po.Object["data"].(map[string]any) + if data["key"] != "v2" { + t.Fatalf("plan object data = %v, want key=v2 (mutated value)", data) + } + + cancel() + if err := <-done; err != nil { + t.Fatalf("collector Run returned %v", err) + } +} diff --git a/test/integration/fakeapiserver/fake.go b/test/integration/fakeapiserver/fake.go new file mode 100644 index 0000000..5167438 --- /dev/null +++ b/test/integration/fakeapiserver/fake.go @@ -0,0 +1,422 @@ +//go:build integration || e2e + +// Package fakeapiserver implements a small in-process Kubernetes API server +// that supports list and watch for one core v1 resource (configmaps) so the +// real watch collector can run against it over the network. +package fakeapiserver + +import ( + "encoding/json" + "fmt" + "net/http" + "net/http/httptest" + "sort" + "strconv" + "strings" + "sync" + + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" +) + +type objState struct { + name string + ns string + obj map[string]any +} + +type eventEntry struct { + Type string + RV string + Obj map[string]any +} + +// Server is a fake Kubernetes apiserver. It is safe for concurrent use. +type Server struct { + URL string + HTTPServer *httptest.Server + + mu sync.Mutex + objects map[string]*objState + events []eventEntry + rvCounter int + forceGap bool + requiredToken string + bookmarkEvery int + broadcast chan struct{} +} + +// NewFake builds a fake apiserver seeded with the given configmaps +// (name to JSON body). The initial resource version is "100". +func NewFake(initial map[string]string) (*Server, error) { + s := &Server{ + objects: map[string]*objState{}, + rvCounter: 100, + broadcast: make(chan struct{}), + } + for name, body := range initial { + var obj map[string]any + if err := json.Unmarshal([]byte(body), &obj); err != nil { + return nil, fmt.Errorf("fakeapiserver: initial configmap %q: %w", name, err) + } + ns, obj := normalize(name, obj, "100") + s.objects[ns+"/"+name] = &objState{name: name, ns: ns, obj: obj} + } + + mux := http.NewServeMux() + mux.HandleFunc("GET /api/v1/namespaces/{ns}/configmaps", s.serveNamespaced) + mux.HandleFunc("GET /api/v1/configmaps", s.serveCluster) + mux.HandleFunc("/", func(w http.ResponseWriter, r *http.Request) { + http.NotFound(w, r) + }) + hts := httptest.NewServer(mux) + s.URL = hts.URL + s.HTTPServer = hts + return s, nil +} + +// RequireToken makes the fake reject any request whose bearer token does not +// match the given value with a 403. Requests without an Authorization header +// are always rejected with a 401. +func (s *Server) RequireToken(tok string) { + s.mu.Lock() + defer s.mu.Unlock() + s.requiredToken = tok +} + +// SetBookmarkEvery enables emitting a BOOKMARK after every n delivered watch +// events. Zero disables bookmarks. +func (s *Server) SetBookmarkEvery(n int) { + s.mu.Lock() + defer s.mu.Unlock() + s.bookmarkEvery = n +} + +// ForceGap marks every currently valid resource version as expired and bumps +// the collection resource version, simulating events that happened while a +// client was disconnected. +func (s *Server) ForceGap() { + s.mu.Lock() + defer s.mu.Unlock() + s.forceGap = true + s.rvCounter++ +} + +// CloseCurrentWatch closes all active client connections, forcing the +// collector to reconnect from its last durable resource version. +func (s *Server) CloseCurrentWatch() { + s.HTTPServer.CloseClientConnections() +} + +// Close shuts the fake apiserver down. +func (s *Server) Close() { + s.HTTPServer.Close() +} + +// AddOrUpdate creates or modifies a configmap. An empty rv assigns the next +// resource version from the server counter; a non-empty rv is used directly. +func (s *Server) AddOrUpdate(name string, obj map[string]any, rv string) { + s.mu.Lock() + defer s.mu.Unlock() + if rv == "" { + s.rvCounter++ + rv = strconv.Itoa(s.rvCounter) + } else if n := parseRV(rv); n > s.rvCounter { + s.rvCounter = n + } + ns, obj := normalize(name, obj, rv) + key := ns + "/" + name + typ := "ADDED" + if _, ok := s.objects[key]; ok { + typ = "MODIFIED" + } + s.objects[key] = &objState{name: name, ns: ns, obj: obj} + s.appendEventLocked(typ, rv, obj) +} + +// Delete removes a configmap by name across all namespaces. +func (s *Server) Delete(name string) { + s.mu.Lock() + defer s.mu.Unlock() + for key, o := range s.objects { + if o.name != name { + continue + } + s.rvCounter++ + rv := strconv.Itoa(s.rvCounter) + del := cloneMap(o.obj) + m, _ := del["metadata"].(map[string]any) + if m == nil { + m = map[string]any{} + del["metadata"] = m + } + m["resourceVersion"] = rv + s.appendEventLocked("DELETED", rv, del) + delete(s.objects, key) + return + } +} + +func (s *Server) serveNamespaced(w http.ResponseWriter, r *http.Request) { + s.serveConfigMaps(w, r, r.PathValue("ns")) +} + +func (s *Server) serveCluster(w http.ResponseWriter, r *http.Request) { + s.serveConfigMaps(w, r, "") +} + +func (s *Server) serveConfigMaps(w http.ResponseWriter, r *http.Request, ns string) { + if !s.authorize(w, r) { + return + } + q := r.URL.Query() + if q.Get("watch") == "true" { + s.serveWatch(w, r, ns, q.Get("resourceVersion"), q.Get("allowWatchBookmarks") == "true") + return + } + s.serveList(w, r, ns) +} + +func (s *Server) authorize(w http.ResponseWriter, r *http.Request) bool { + tok := "" + if h := r.Header.Get("Authorization"); h != "" { + tok = strings.TrimPrefix(h, "Bearer ") + } + if tok == "" { + writeStatus(w, http.StatusUnauthorized, metav1.StatusReasonUnauthorized, "missing Authorization header") + return false + } + s.mu.Lock() + required := s.requiredToken + s.mu.Unlock() + if required != "" && tok != required { + writeStatus(w, http.StatusForbidden, metav1.StatusReasonForbidden, "the provided token does not match the required token") + return false + } + return true +} + +func (s *Server) serveList(w http.ResponseWriter, r *http.Request, ns string) { + s.mu.Lock() + if s.forceGap { + if rv := r.URL.Query().Get("resourceVersion"); rv != "" && parseRV(rv) < s.rvCounter { + s.forceGap = false + s.mu.Unlock() + writeStatus(w, http.StatusGone, metav1.StatusReasonExpired, "resourceVersion="+rv+" is expired") + return + } + } + current := strconv.Itoa(s.rvCounter) + items := make([]map[string]any, 0, len(s.objects)) + for _, o := range s.objects { + if ns != "" && o.ns != ns { + continue + } + items = append(items, o.obj) + } + s.mu.Unlock() + sort.Slice(items, func(i, j int) bool { + ni, nj := objNS(items[i]), objNS(items[j]) + if ni != nj { + return ni < nj + } + return objName(items[i]) < objName(items[j]) + }) + w.Header().Set("Content-Type", "application/json") + json.NewEncoder(w).Encode(map[string]any{ + "apiVersion": "v1", + "kind": "ConfigMapList", + "metadata": map[string]any{"resourceVersion": current}, + "items": items, + }) +} + +func (s *Server) serveWatch(w http.ResponseWriter, r *http.Request, ns, rv string, bookmarks bool) { + s.mu.Lock() + if s.forceGap && rv != "" && parseRV(rv) < s.rvCounter { + s.forceGap = false + s.mu.Unlock() + s.streamGone(w, "resourceVersion="+rv+" is expired") + return + } + every := s.bookmarkEvery + s.mu.Unlock() + + fl, ok := w.(http.Flusher) + if !ok { + http.Error(w, "streaming unsupported", http.StatusInternalServerError) + return + } + w.Header().Set("Content-Type", "application/json") + w.WriteHeader(http.StatusOK) + fl.Flush() + + s.mu.Lock() + startIdx := len(s.events) + if rv != "" { + if i := s.firstEventAfterLocked(rv); i >= 0 { + startIdx = i + } + } + s.mu.Unlock() + + idx := startIdx + delivered := 0 + for { + s.mu.Lock() + if idx < len(s.events) { + ev := s.events[idx] + idx++ + s.mu.Unlock() + if err := writeJSONLine(w, fl, map[string]any{"type": ev.Type, "object": ev.Obj}); err != nil { + return + } + delivered++ + if bookmarks && every > 0 && delivered%every == 0 { + if err := s.writeBookmark(w, fl); err != nil { + return + } + } + continue + } + ch := s.broadcast + s.mu.Unlock() + select { + case <-ch: + case <-r.Context().Done(): + return + } + } +} + +func (s *Server) writeBookmark(w http.ResponseWriter, fl http.Flusher) error { + s.mu.Lock() + rv := strconv.Itoa(s.rvCounter) + s.mu.Unlock() + return writeJSONLine(w, fl, map[string]any{ + "type": "BOOKMARK", + "object": map[string]any{ + "kind": "ConfigMap", + "apiVersion": "v1", + "metadata": map[string]any{"resourceVersion": rv}, + }, + }) +} + +func (s *Server) streamGone(w http.ResponseWriter, msg string) { + fl, ok := w.(http.Flusher) + if !ok { + http.Error(w, "streaming unsupported", http.StatusInternalServerError) + return + } + w.Header().Set("Content-Type", "application/json") + w.WriteHeader(http.StatusOK) + fl.Flush() + writeJSONLine(w, fl, map[string]any{ + "type": "ERROR", + "object": metav1.Status{ + TypeMeta: metav1.TypeMeta{Kind: "Status", APIVersion: "v1"}, + Status: metav1.StatusFailure, + Message: msg, + Reason: metav1.StatusReasonExpired, + Code: http.StatusGone, + }, + }) +} + +func (s *Server) appendEventLocked(typ, rv string, obj map[string]any) { + s.events = append(s.events, eventEntry{Type: typ, RV: rv, Obj: obj}) + close(s.broadcast) + s.broadcast = make(chan struct{}) +} + +func (s *Server) firstEventAfterLocked(rv string) int { + want := parseRV(rv) + for i := range s.events { + if parseRV(s.events[i].RV) > want { + return i + } + } + return len(s.events) +} + +func normalize(name string, obj map[string]any, rv string) (string, map[string]any) { + ns := "default" + m, _ := obj["metadata"].(map[string]any) + if m == nil { + m = map[string]any{} + obj["metadata"] = m + } + if n, _ := m["namespace"].(string); n != "" { + ns = n + } + m["name"] = name + m["namespace"] = ns + m["resourceVersion"] = rv + if _, ok := m["uid"]; !ok { + m["uid"] = "uid-" + name + } + if _, ok := obj["kind"]; !ok { + obj["kind"] = "ConfigMap" + } + if _, ok := obj["apiVersion"]; !ok { + obj["apiVersion"] = "v1" + } + return ns, obj +} + +func writeStatus(w http.ResponseWriter, code int32, reason metav1.StatusReason, msg string) { + w.Header().Set("Content-Type", "application/json") + w.WriteHeader(int(code)) + json.NewEncoder(w).Encode(metav1.Status{ + TypeMeta: metav1.TypeMeta{Kind: "Status", APIVersion: "v1"}, + Status: metav1.StatusFailure, + Message: msg, + Reason: reason, + Code: code, + }) +} + +func writeJSONLine(w http.ResponseWriter, fl http.Flusher, v any) error { + data, err := json.Marshal(v) + if err != nil { + return err + } + if _, err := w.Write(append(data, '\n')); err != nil { + return err + } + fl.Flush() + return nil +} + +func parseRV(rv string) int { + if rv == "" { + return 0 + } + n, err := strconv.Atoi(rv) + if err != nil { + return 0 + } + return n +} + +func objName(obj map[string]any) string { + m, _ := obj["metadata"].(map[string]any) + name, _ := m["name"].(string) + return name +} + +func objNS(obj map[string]any) string { + m, _ := obj["metadata"].(map[string]any) + ns, _ := m["namespace"].(string) + return ns +} + +func cloneMap(m map[string]any) map[string]any { + b, err := json.Marshal(m) + if err != nil { + return map[string]any{} + } + var out map[string]any + _ = json.Unmarshal(b, &out) + return out +} diff --git a/test/integration/fakeapiserver/fake_test.go b/test/integration/fakeapiserver/fake_test.go new file mode 100644 index 0000000..0cb42e4 --- /dev/null +++ b/test/integration/fakeapiserver/fake_test.go @@ -0,0 +1,301 @@ +//go:build integration || e2e + +package fakeapiserver + +import ( + "bufio" + "encoding/json" + "io" + "net/http" + "testing" +) + +type watchEvent struct { + Type string `json:"type"` + Object json.RawMessage `json:"object"` +} + +func doRaw(t *testing.T, method, url string) *http.Response { + t.Helper() + req, err := http.NewRequest(method, url, nil) + if err != nil { + t.Fatalf("new request: %v", err) + } + req.Header.Set("Authorization", "Bearer test-token") + resp, err := http.DefaultClient.Do(req) + if err != nil { + t.Fatalf("%s %s: %v", method, url, err) + } + return resp +} + +func startWatch(t *testing.T, s *Server, rv string, bookmarks bool) (*http.Response, *bufio.Reader) { + t.Helper() + url := s.URL + "/api/v1/namespaces/default/configmaps?watch=true" + if rv != "" { + url += "&resourceVersion=" + rv + } + if bookmarks { + url += "&allowWatchBookmarks=true" + } + resp := doRaw(t, http.MethodGet, url) + if resp.StatusCode != http.StatusOK { + body, _ := io.ReadAll(resp.Body) + t.Fatalf("watch start status = %d, body = %s", resp.StatusCode, body) + } + return resp, bufio.NewReader(resp.Body) +} + +func nextEvent(t *testing.T, r *bufio.Reader) watchEvent { + t.Helper() + line, err := r.ReadString('\n') + if err != nil { + t.Fatalf("read watch line: %v", err) + } + var ev watchEvent + if err := json.Unmarshal([]byte(line), &ev); err != nil { + t.Fatalf("decode watch line %q: %v", line, err) + } + return ev +} + +func cm(name string) map[string]any { + return map[string]any{ + "metadata": map[string]any{"name": name}, + "data": map[string]any{"app": name}, + } +} + +func TestListReturnsConfigMapList(t *testing.T) { + initial := map[string]string{ + "cm-b": `{"metadata":{"name":"cm-b"},"data":{"k":"v"}}`, + "cm-a": `{"metadata":{"name":"cm-a"},"data":{"k":"v"}}`, + } + s, err := NewFake(initial) + if err != nil { + t.Fatalf("NewFake: %v", err) + } + defer s.Close() + + resp := doRaw(t, http.MethodGet, s.URL+"/api/v1/namespaces/default/configmaps") + defer resp.Body.Close() + if resp.StatusCode != http.StatusOK { + body, _ := io.ReadAll(resp.Body) + t.Fatalf("list status = %d, body = %s", resp.StatusCode, body) + } + + var list struct { + APIVersion string `json:"apiVersion"` + Kind string `json:"kind"` + Metadata struct { + ResourceVersion string `json:"resourceVersion"` + } `json:"metadata"` + Items []map[string]any `json:"items"` + } + if err := json.NewDecoder(resp.Body).Decode(&list); err != nil { + t.Fatalf("decode list: %v", err) + } + if list.APIVersion != "v1" || list.Kind != "ConfigMapList" { + t.Fatalf("list apiVersion/kind = %q/%q", list.APIVersion, list.Kind) + } + if list.Metadata.ResourceVersion != "100" { + t.Fatalf("list resourceVersion = %q, want 100", list.Metadata.ResourceVersion) + } + if len(list.Items) != 2 { + t.Fatalf("list items = %d, want 2", len(list.Items)) + } + if got := objName(list.Items[0]); got != "cm-a" { + t.Fatalf("first item = %q, want cm-a (sorted)", got) + } +} + +func TestWatchReceivesAddedAfterAdd(t *testing.T) { + s, err := NewFake(nil) + if err != nil { + t.Fatalf("NewFake: %v", err) + } + defer s.Close() + + resp, r := startWatch(t, s, "100", false) + defer resp.Body.Close() + + s.AddOrUpdate("cm-a", cm("cm-a"), "") + + ev := nextEvent(t, r) + if ev.Type != "ADDED" { + t.Fatalf("first event type = %q, want ADDED", ev.Type) + } + var obj map[string]any + if err := json.Unmarshal(ev.Object, &obj); err != nil { + t.Fatalf("decode event object: %v", err) + } + if objName(obj) != "cm-a" { + t.Fatalf("event object name = %q, want cm-a", objName(obj)) + } + if m, _ := obj["metadata"].(map[string]any); m["resourceVersion"] != "101" { + t.Fatalf("event object resourceVersion = %v, want 101", m["resourceVersion"]) + } +} + +func TestWatchDeliversBookmark(t *testing.T) { + s, err := NewFake(nil) + if err != nil { + t.Fatalf("NewFake: %v", err) + } + defer s.Close() + s.SetBookmarkEvery(1) + + resp, r := startWatch(t, s, "100", true) + defer resp.Body.Close() + + s.AddOrUpdate("cm-a", cm("cm-a"), "") + + if ev := nextEvent(t, r); ev.Type != "ADDED" { + t.Fatalf("first event type = %q, want ADDED", ev.Type) + } + ev := nextEvent(t, r) + if ev.Type != "BOOKMARK" { + t.Fatalf("second event type = %q, want BOOKMARK", ev.Type) + } + var obj map[string]any + if err := json.Unmarshal(ev.Object, &obj); err != nil { + t.Fatalf("decode bookmark object: %v", err) + } + m, _ := obj["metadata"].(map[string]any) + if m["resourceVersion"] == "" { + t.Fatal("bookmark missing resourceVersion") + } +} + +func TestForceGapExpiresOldResourceVersion(t *testing.T) { + s, err := NewFake(nil) + if err != nil { + t.Fatalf("NewFake: %v", err) + } + defer s.Close() + + s.AddOrUpdate("cm-a", cm("cm-a"), "") + s.ForceGap() + + resp, r := startWatch(t, s, "100", false) + defer resp.Body.Close() + + ev := nextEvent(t, r) + if ev.Type != "ERROR" { + t.Fatalf("event type = %q, want ERROR", ev.Type) + } + var status struct { + Kind string `json:"kind"` + Status string `json:"status"` + Reason string `json:"reason"` + Code int `json:"code"` + } + if err := json.Unmarshal(ev.Object, &status); err != nil { + t.Fatalf("decode status object: %v", err) + } + if status.Kind != "Status" || status.Status != "Failure" || status.Reason != "Expired" || status.Code != http.StatusGone { + t.Fatalf("status object = %+v, want Status/Failure/Expired/410", status) + } +} + +func TestForceGapLeavesCurrentResourceVersionWatchable(t *testing.T) { + s, err := NewFake(nil) + if err != nil { + t.Fatalf("NewFake: %v", err) + } + defer s.Close() + + s.AddOrUpdate("cm-a", cm("cm-a"), "") + s.ForceGap() + + resp, r := startWatch(t, s, "102", false) + defer resp.Body.Close() + + s.AddOrUpdate("cm-b", cm("cm-b"), "") + ev := nextEvent(t, r) + if ev.Type != "ADDED" { + t.Fatalf("event type = %q, want ADDED", ev.Type) + } + var obj map[string]any + if err := json.Unmarshal(ev.Object, &obj); err != nil { + t.Fatalf("decode event object: %v", err) + } + if objName(obj) != "cm-b" { + t.Fatalf("event object name = %q, want cm-b", objName(obj)) + } +} + +func TestWatchRequiresAuthorizationHeader(t *testing.T) { + s, err := NewFake(nil) + if err != nil { + t.Fatalf("NewFake: %v", err) + } + defer s.Close() + + for _, path := range []string{ + "/api/v1/namespaces/default/configmaps", + "/api/v1/configmaps", + "/api/v1/namespaces/default/configmaps?watch=true&resourceVersion=100", + } { + req, err := http.NewRequest(http.MethodGet, s.URL+path, nil) + if err != nil { + t.Fatalf("new request: %v", err) + } + resp, err := http.DefaultClient.Do(req) + if err != nil { + t.Fatalf("%s: %v", path, err) + } + body, _ := io.ReadAll(resp.Body) + resp.Body.Close() + if resp.StatusCode != http.StatusUnauthorized { + t.Fatalf("%s status = %d, want 401; body = %s", path, resp.StatusCode, body) + } + } +} + +func TestWatchRejectsWrongRequiredToken(t *testing.T) { + s, err := NewFake(nil) + if err != nil { + t.Fatalf("NewFake: %v", err) + } + defer s.Close() + s.RequireToken("super-secret") + + req, _ := http.NewRequest(http.MethodGet, s.URL+"/api/v1/configmaps", nil) + req.Header.Set("Authorization", "Bearer wrong-token") + resp, err := http.DefaultClient.Do(req) + if err != nil { + t.Fatalf("request: %v", err) + } + defer resp.Body.Close() + if resp.StatusCode != http.StatusForbidden { + body, _ := io.ReadAll(resp.Body) + t.Fatalf("status = %d, want 403; body = %s", resp.StatusCode, body) + } +} + +func TestClusterWideListIncludesAllNamespaces(t *testing.T) { + s, err := NewFake(nil) + if err != nil { + t.Fatalf("NewFake: %v", err) + } + defer s.Close() + + other := map[string]any{ + "metadata": map[string]any{"name": "cm-other", "namespace": "kube-system"}, + } + s.AddOrUpdate("cm-other", other, "") + s.AddOrUpdate("cm-a", cm("cm-a"), "") + + resp := doRaw(t, http.MethodGet, s.URL+"/api/v1/configmaps") + defer resp.Body.Close() + var list struct { + Items []map[string]any `json:"items"` + } + if err := json.NewDecoder(resp.Body).Decode(&list); err != nil { + t.Fatalf("decode list: %v", err) + } + if len(list.Items) != 2 { + t.Fatalf("cluster-wide items = %d, want 2", len(list.Items)) + } +} diff --git a/test/integration/watch_integration_test.go b/test/integration/watch_integration_test.go new file mode 100644 index 0000000..b22a15d --- /dev/null +++ b/test/integration/watch_integration_test.go @@ -0,0 +1,331 @@ +//go:build integration + +// Package integration_test runs the real watch collector and storage against +// a fake Kubernetes apiserver over the network. +package integration_test + +import ( + "context" + "encoding/json" + "testing" + "time" + + "k8s.io/apimachinery/pkg/runtime/schema" + "k8s.io/client-go/dynamic" + "k8s.io/client-go/kubernetes/scheme" + "k8s.io/client-go/rest" + + "github.com/krply/krply/internal/discovery" + "github.com/krply/krply/internal/event" + "github.com/krply/krply/internal/storage" + "github.com/krply/krply/internal/watch" + "github.com/krply/krply/test/integration/fakeapiserver" +) + +const cmA = `{"metadata":{"name":"cm-a"},"data":{"key":"v1"}}` + +func newDynamicClient(t *testing.T, url, token string) dynamic.Interface { + t.Helper() + cfg := &rest.Config{ + Host: url, + ContentConfig: rest.ContentConfig{ + NegotiatedSerializer: scheme.Codecs.WithoutConversion(), + GroupVersion: &schema.GroupVersion{Group: "", Version: "v1"}, + }, + } + if token != "" { + cfg.BearerToken = token + } + dyn, err := dynamic.NewForConfig(cfg) + if err != nil { + t.Fatalf("dynamic.NewForConfig: %v", err) + } + return dyn +} + +func newCollector(t *testing.T, dyn dynamic.Interface, store storage.Store) *watch.Collector { + t.Helper() + col, err := watch.NewCollector(watch.Config{ + ClusterID: "itest", + Resources: []discovery.ResourceSpec{ + {APIGroup: "", Version: "v1", Resource: "configmaps", Kind: "ConfigMap", Namespace: "default"}, + }, + Store: store, + DynamicClient: dyn, + Bookmarks: true, + MinBackoff: 10 * time.Millisecond, + MaxBackoff: 100 * time.Millisecond, + }) + if err != nil { + t.Fatalf("watch.NewCollector: %v", err) + } + return col +} + +func waitFor(t *testing.T, timeout time.Duration, cond func() bool, msg string) { + t.Helper() + deadline := time.Now().Add(timeout) + for time.Now().Before(deadline) { + if cond() { + return + } + time.Sleep(10 * time.Millisecond) + } + t.Fatalf("timed out waiting for %s", msg) +} + +func storeEvents(ctx context.Context, store storage.Store) []event.Record { + recs, err := store.Events(ctx, storage.EventFilter{}) + if err != nil { + return nil + } + return recs +} + +func baselineAndSynthetic(store storage.Store, name string) func() bool { + return func() bool { + var baseline, synthetic bool + for _, r := range storeEvents(context.Background(), store) { + switch { + case r.Type == event.TypeBaseline: + baseline = true + case r.Type == event.TypeEvent && r.Synthetic && r.WatchType == event.WatchAdded && r.Resource.Name == name: + synthetic = true + } + } + return baseline && synthetic + } +} + +// TestCollectorRecordsBaselineAndModified verifies the initial list produces a +// BASELINE plus a synthetic ADDED, and that a mutation through the fake +// apiserver reaches the store as a single MODIFIED event with a stable +// EventID. +func TestCollectorRecordsBaselineAndModified(t *testing.T) { + fake, err := fakeapiserver.NewFake(map[string]string{"cm-a": cmA}) + if err != nil { + t.Fatalf("NewFake: %v", err) + } + defer fake.Close() + + store := storage.NewInMemory() + col := newCollector(t, newDynamicClient(t, fake.URL, "ignored"), store) + + ctx, cancel := context.WithCancel(context.Background()) + done := make(chan error, 1) + go func() { done <- col.Run(ctx) }() + + waitFor(t, 5*time.Second, baselineAndSynthetic(store, "cm-a"), "baseline + synthetic ADDED for cm-a") + + fake.AddOrUpdate("cm-a", map[string]any{ + "metadata": map[string]any{"name": "cm-a"}, + "data": map[string]any{"key": "v2"}, + }, "") + + var modified *event.Record + waitFor(t, 5*time.Second, func() bool { + for _, r := range storeEvents(ctx, store) { + if r.Type == event.TypeEvent && r.WatchType == event.WatchModified && r.Resource.Name == "cm-a" { + modified = &r + return true + } + } + return false + }, "MODIFIED event for cm-a") + + if modified.WatchType != event.WatchModified { + t.Fatalf("watch type = %q, want MODIFIED", modified.WatchType) + } + var obj map[string]any + if err := json.Unmarshal(modified.Object, &obj); err != nil { + t.Fatalf("decode modified object: %v", err) + } + data, _ := obj["data"].(map[string]any) + if data["key"] != "v2" { + t.Fatalf("modified data = %v, want key=v2", data) + } + if modified.EventID == "" { + t.Fatal("MODIFIED event missing EventID") + } + + count := 0 + for _, r := range storeEvents(ctx, store) { + if r.EventID == modified.EventID { + count++ + } + } + if count != 1 { + t.Fatalf("EventID %s appeared %d times, want exactly once", modified.EventID, count) + } + + cancel() + if err := <-done; err != nil { + t.Fatalf("collector Run returned %v", err) + } +} + +// TestCollectorRecordsCheckpoint verifies bookmarks from the fake apiserver +// advance checkpoints only and are not stored as events with object payloads. +func TestCollectorRecordsCheckpoint(t *testing.T) { + fake, err := fakeapiserver.NewFake(map[string]string{"cm-a": cmA}) + if err != nil { + t.Fatalf("NewFake: %v", err) + } + defer fake.Close() + fake.SetBookmarkEvery(1) + + store := storage.NewInMemory() + col := newCollector(t, newDynamicClient(t, fake.URL, "ignored"), store) + + ctx, cancel := context.WithCancel(context.Background()) + done := make(chan error, 1) + go func() { done <- col.Run(ctx) }() + defer cancel() + + waitFor(t, 5*time.Second, baselineAndSynthetic(store, "cm-a"), "baseline + synthetic ADDED for cm-a") + + fake.AddOrUpdate("cm-a", map[string]any{ + "metadata": map[string]any{"name": "cm-a"}, + "data": map[string]any{"key": "v3"}, + }, "") + + var checkpoint *event.Record + waitFor(t, 5*time.Second, func() bool { + for _, r := range storeEvents(ctx, store) { + if r.Type == event.TypeCheckpoint && r.Checkpoint != nil && r.Checkpoint.ResourceVersion != "" { + checkpoint = &r + return true + } + } + return false + }, "CHECKPOINT record") + + if checkpoint.Checkpoint == nil || checkpoint.Checkpoint.ResourceVersion == "" { + t.Fatal("CHECKPOINT missing resource version") + } + if len(checkpoint.Object) != 0 { + t.Fatalf("CHECKPOINT must not carry an object payload, got %d bytes", len(checkpoint.Object)) + } + if checkpoint.WatchType != "" { + t.Fatalf("CHECKPOINT must not be an EVENT, got watch type %q", checkpoint.WatchType) + } + + select { + case err := <-done: + t.Fatalf("collector Run returned early with %v", err) + default: + } + cancel() + if err := <-done; err != nil { + t.Fatalf("collector Run returned %v", err) + } +} + +// TestCollectorRelistsAfterGap verifies that a forced 410 on a reconnect +// writes a GAP record and triggers a relist producing a second BASELINE. +func TestCollectorRelistsAfterGap(t *testing.T) { + fake, err := fakeapiserver.NewFake(map[string]string{"cm-a": cmA}) + if err != nil { + t.Fatalf("NewFake: %v", err) + } + defer fake.Close() + + store := storage.NewInMemory() + col := newCollector(t, newDynamicClient(t, fake.URL, "ignored"), store) + + ctx, cancel := context.WithCancel(context.Background()) + done := make(chan error, 1) + go func() { done <- col.Run(ctx) }() + defer cancel() + + waitFor(t, 5*time.Second, baselineAndSynthetic(store, "cm-a"), "baseline + synthetic ADDED for cm-a") + + fake.ForceGap() + fake.CloseCurrentWatch() + + waitFor(t, 5*time.Second, func() bool { + for _, r := range storeEvents(ctx, store) { + if r.Type == event.TypeGap && r.Gap != nil && r.Gap.Reason == "410 Gone" { + return true + } + } + return false + }, "GAP record with reason 410 Gone") + + waitFor(t, 5*time.Second, func() bool { + baselines := 0 + for _, r := range storeEvents(ctx, store) { + if r.Type == event.TypeBaseline { + baselines++ + } + } + return baselines >= 2 + }, "relist after gap (second BASELINE)") + + waitFor(t, 5*time.Second, func() bool { + synthetic := 0 + for _, r := range storeEvents(ctx, store) { + if r.Type == event.TypeEvent && r.Synthetic && r.WatchType == event.WatchAdded && r.Resource.Name == "cm-a" { + synthetic++ + } + } + return synthetic >= 2 + }, "second synthetic ADDED for cm-a after relist") + + gaps := 0 + baselines := 0 + for _, r := range storeEvents(ctx, store) { + switch r.Type { + case event.TypeGap: + gaps++ + case event.TypeBaseline: + baselines++ + } + } + if gaps != 1 { + t.Fatalf("GAP records = %d, want exactly 1", gaps) + } + if baselines != 2 { + t.Fatalf("BASELINE records = %d, want exactly 2", baselines) + } +} + +// TestCollectorRetriesOnRbacDenial verifies that a collector whose dynamic +// client is denied by the apiserver keeps retrying instead of crashing or +// exiting, and stops only when the context is cancelled. +func TestCollectorRetriesOnRbacDenial(t *testing.T) { + fake, err := fakeapiserver.NewFake(nil) + if err != nil { + t.Fatalf("NewFake: %v", err) + } + defer fake.Close() + fake.RequireToken("super-secret") + + store := storage.NewInMemory() + col := newCollector(t, newDynamicClient(t, fake.URL, "wrong-token"), store) + + ctx, cancel := context.WithCancel(context.Background()) + done := make(chan error, 1) + go func() { done <- col.Run(ctx) }() + + time.Sleep(200 * time.Millisecond) + select { + case err := <-done: + t.Fatalf("collector Run returned early (err = %v); RBAC denial should retry", err) + default: + } + + cancel() + select { + case err := <-done: + if err != nil { + t.Fatalf("collector Run returned %v after cancel, want nil", err) + } + case <-time.After(5 * time.Second): + t.Fatal("collector Run did not return after cancel") + } + + if recs := storeEvents(ctx, store); len(recs) != 0 { + t.Fatalf("RBAC-denied collector wrote %d records, want none", len(recs)) + } +}