test: add unit, integration, and e2e suites

- unit tests across storage, replay, and api layers
- integration suite driven by a fake apiserver (no cluster needed)
- end-to-end suite against a live cluster, tagged e2e
This commit is contained in:
lakshit verma 2026-08-06 05:34:04 +05:30
parent b04fd7a406
commit 3eb2a2ffd4
No known key found for this signature in database
GPG key ID: EB498AFC60A7A01A
4 changed files with 1279 additions and 0 deletions

225
test/e2e/e2e_test.go Normal file
View file

@ -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)
}
}

View file

@ -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
}

View file

@ -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))
}
}

View file

@ -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))
}
}