From 45c8daa1f55eee2edd59e265c10cdf421104b901 Mon Sep 17 00:00:00 2001 From: Geoffrey Ragot Date: Tue, 2 Jun 2026 08:25:57 +0200 Subject: [PATCH 1/3] feat: use Server-Side Apply for CRD management Replace MergePatch with Server-Side Apply (SSA) so the agent only owns the fields it explicitly sets. This allows users to manually edit CRs without the agent overwriting their changes on resync. - Field manager "formance-agent" for main sync (force=false, conflicts are logged and skipped) - Field manager "formance-agent:disable" for disable/enable commands (force=true, authoritative) - Add integration tests verifying user annotations survive resync Co-Authored-By: Claude Opus 4.6 (1M context) --- internal/k8s_client.go | 35 ++++----- internal/membership_listener.go | 113 +++++++++++++++++------------- tests/membership_listener_test.go | 66 +++++++++++++++++ 3 files changed, 151 insertions(+), 63 deletions(-) diff --git a/internal/k8s_client.go b/internal/k8s_client.go index 3537201..1d1e763 100644 --- a/internal/k8s_client.go +++ b/internal/k8s_client.go @@ -2,6 +2,8 @@ package internal import ( "context" + "encoding/json" + "fmt" "github.com/formancehq/go-libs/v2/collectionutils" "github.com/formancehq/go-libs/v2/logging" @@ -19,8 +21,7 @@ import ( type K8SClient interface { Get(ctx context.Context, resource string, name string) (*unstructured.Unstructured, error) - Create(ctx context.Context, resource string, o *unstructured.Unstructured) error - Patch(ctx context.Context, resource, name string, body []byte) error + Apply(ctx context.Context, resource string, obj *unstructured.Unstructured, fieldManager string, force bool) (*unstructured.Unstructured, error) Delete(ctx context.Context, resource, name string) error EnsureNotExists(ctx context.Context, resource, name string) error EnsureNotExistsBySelector(ctx context.Context, resource string, selector labels.Selector) error @@ -43,22 +44,24 @@ func (c defaultK8SClient) Get(ctx context.Context, resource string, name string) return u, nil } -func (c defaultK8SClient) Create(ctx context.Context, resource string, o *unstructured.Unstructured) error { - return c.restClient. - Post(). - Resource(resource). - Body(o). - Do(ctx). - Into(o) -} - -func (c defaultK8SClient) Patch(ctx context.Context, resource, name string, body []byte) error { - return c.restClient.Patch(types.MergePatchType). - Name(name). - Body(body). +func (c defaultK8SClient) Apply(ctx context.Context, resource string, obj *unstructured.Unstructured, fieldManager string, force bool) (*unstructured.Unstructured, error) { + data, err := json.Marshal(obj.Object) + if err != nil { + return nil, err + } + result := &unstructured.Unstructured{} + err = c.restClient.Patch(types.ApplyPatchType). Resource(resource). + Name(obj.GetName()). + Param("fieldManager", fieldManager). + Param("force", fmt.Sprintf("%t", force)). + Body(data). Do(ctx). - Error() + Into(result) + if err != nil { + return nil, err + } + return result, nil } func (c defaultK8SClient) Delete(ctx context.Context, resource, name string) error { diff --git a/internal/membership_listener.go b/internal/membership_listener.go index b67a450..f82def1 100644 --- a/internal/membership_listener.go +++ b/internal/membership_listener.go @@ -3,7 +3,6 @@ package internal import ( "context" - "encoding/json" "fmt" "net/url" "slices" @@ -20,7 +19,6 @@ import ( "github.com/formancehq/stack/components/agent/internal/generated" "github.com/formancehq/stack/components/agent/internal/grpcclient" "github.com/pkg/errors" - "k8s.io/apimachinery/pkg/api/equality" apierrors "k8s.io/apimachinery/pkg/api/errors" "k8s.io/apimachinery/pkg/api/meta" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" @@ -29,6 +27,11 @@ import ( "sigs.k8s.io/controller-runtime/pkg/client" ) +const ( + fieldManagerAgent = "formance-agent" + fieldManagerDisable = "formance-agent:disable" +) + //go:generate mockgen -source=membership_listener.go -destination=membership_client_generated.go -package=internal . MembershipClient type MembershipClient interface { Orders() chan *generated.Order @@ -363,7 +366,19 @@ func (c *membershipListener) deleteStack(ctx context.Context, stack *generated.D } func (c *membershipListener) disableStack(ctx context.Context, stack *generated.DisabledStack) { - if err := c.client.Patch(ctx, "Stacks", stack.ClusterName, []byte(`{"spec": {"disabled": true}}`)); err != nil { + obj := &unstructured.Unstructured{ + Object: map[string]any{ + "apiVersion": formanceGroupVersion.String(), + "kind": "Stack", + "metadata": map[string]any{ + "name": stack.ClusterName, + }, + "spec": map[string]any{ + "disabled": true, + }, + }, + } + if _, err := c.client.Apply(ctx, "Stacks", obj, fieldManagerDisable, true); err != nil { logging.FromContext(ctx).Errorf("Disabling cluster side: %s", err) return } @@ -372,8 +387,20 @@ func (c *membershipListener) disableStack(ctx context.Context, stack *generated. } func (c *membershipListener) enableStack(ctx context.Context, stack *generated.EnabledStack) { - if err := c.client.Patch(ctx, "Stacks", stack.ClusterName, []byte(`{"spec": {"disabled": false}}`)); err != nil { - logging.FromContext(ctx).Errorf("Disabling cluster side: %s", err) + obj := &unstructured.Unstructured{ + Object: map[string]any{ + "apiVersion": formanceGroupVersion.String(), + "kind": "Stack", + "metadata": map[string]any{ + "name": stack.ClusterName, + }, + "spec": map[string]any{ + "disabled": false, + }, + }, + } + if _, err := c.client.Apply(ctx, "Stacks", obj, fieldManagerDisable, true); err != nil { + logging.FromContext(ctx).Errorf("Enabling cluster side: %s", err) return } @@ -385,64 +412,56 @@ func (c *membershipListener) createOrUpdate(ctx context.Context, gvk schema.Grou logger := logging.FromContext(ctx).WithFields(map[string]any{ "gvk": gvk, }) - logger.Infof("creating object '%s'", name) + logger.Infof("applying object '%s'", name) + if content["metadata"] == nil { content["metadata"] = map[string]any{} } + md := content["metadata"].(map[string]any) - if content["metadata"].(map[string]any)["labels"] == nil { - content["metadata"].(map[string]any)["labels"] = map[string]any{} + if md["labels"] == nil { + md["labels"] = map[string]any{} + } + md["labels"].(map[string]any)["formance.com/created-by-agent"] = "true" + md["labels"].(map[string]any)["formance.com/stack"] = stackName + md["name"] = name + + if owner != nil { + md["ownerReferences"] = []any{ + map[string]any{ + "apiVersion": owner.APIVersion, + "kind": owner.Kind, + "name": owner.Name, + "uid": string(owner.UID), + }, + } } - content["metadata"].(map[string]any)["labels"].(map[string]any)["formance.com/created-by-agent"] = "true" - content["metadata"].(map[string]any)["labels"].(map[string]any)["formance.com/stack"] = stackName - content["metadata"].(map[string]any)["name"] = name + content["apiVersion"] = gvk.GroupVersion().String() + content["kind"] = gvk.Kind restMapping, err := c.restMapper.RESTMapping(gvk.GroupKind()) if err != nil { return nil, errors.Wrap(err, "getting rest mapping") } - u, err := c.client.Get(ctx, restMapping.Resource.Resource, name) - if err != nil { - if !apierrors.IsNotFound(err) { - return nil, errors.Wrap(err, "reading object") - } - - logger.Infof("Object not found, create a new one") - - u := &unstructured.Unstructured{} - u.SetUnstructuredContent(content) - u.SetGroupVersionKind(gvk) - u.SetName(name) - if owner != nil { - u.SetOwnerReferences([]metav1.OwnerReference{*owner}) - } - - if err := c.client.Create(ctx, restMapping.Resource.Resource, u); err != nil { - return nil, errors.Wrap(err, "creating object") - } - - return u, nil + u := &unstructured.Unstructured{} + u.SetUnstructuredContent(content) - } - - if equality.Semantic.DeepDerivative(content, u.Object) { - logger.Infof("Object found and has expected content, skip it") - return u, nil - } - - logger.Infof("Object exists and content differ, patch it") - contentData, err := json.Marshal(content) + result, err := c.client.Apply(ctx, restMapping.Resource.Resource, u, fieldManagerAgent, false) if err != nil { - return nil, err - } - - if err := c.client.Patch(ctx, restMapping.Resource.Resource, name, contentData); err != nil { - return nil, errors.Wrap(err, "patching object") + if apierrors.IsConflict(err) { + logger.Infof("Conflict applying %s/%s: field owned by another manager, skipping update", gvk.Kind, name) + existing, getErr := c.client.Get(ctx, restMapping.Resource.Resource, name) + if getErr != nil { + return nil, errors.Wrap(getErr, "getting existing object after conflict") + } + return existing, nil + } + return nil, errors.Wrap(err, "applying object") } - return u, nil + return result, nil } func (c *membershipListener) createOrUpdateStackDependency( diff --git a/tests/membership_listener_test.go b/tests/membership_listener_test.go index b2f48f9..caacffa 100644 --- a/tests/membership_listener_test.go +++ b/tests/membership_listener_test.go @@ -227,6 +227,72 @@ var _ = Describe("Membership listener", func() { Expect(u).To(TargetStack(stack)) } }) + When("a user manually edits a resource managed by the agent", func() { + It("Should preserve user-added annotations after resync", func() { + userPatch, err := json.Marshal(map[string]any{ + "metadata": map[string]any{ + "annotations": map[string]any{ + "user-custom/annotation": "user-value", + }, + }, + }) + Expect(err).To(BeNil()) + Expect(k8sClient.Patch(types.MergePatchType). + Resource("Stacks"). + Name(membershipStack.ClusterName). + Body(userPatch). + Do(context.Background()). + Error()).To(Succeed()) + + membershipClient.Orders() <- &generated.Order{ + Message: &generated.Order_ExistingStack{ + ExistingStack: membershipStack, + }, + } + + Eventually(func(g Gomega) { + s := &unstructured.Unstructured{} + g.Expect(LoadResource("Stacks", membershipStack.ClusterName, s)).To(Succeed()) + g.Expect(s.GetAnnotations()).To(HaveKeyWithValue("user-custom/annotation", "user-value")) + versionsFromFile, _, _ := unstructured.NestedString(s.Object, "spec", "versionsFromFile") + g.Expect(versionsFromFile).To(Equal("default")) + }).Should(Succeed()) + }) + + It("Should preserve user-added annotations on modules after resync", func() { + auth := &unstructured.Unstructured{} + Eventually(func() error { + return LoadResource("Auths", membershipStack.ClusterName, auth) + }).Should(BeNil()) + + userPatch, err := json.Marshal(map[string]any{ + "metadata": map[string]any{ + "annotations": map[string]any{ + "user-custom/note": "do not delete", + }, + }, + }) + Expect(err).To(BeNil()) + Expect(k8sClient.Patch(types.MergePatchType). + Resource("Auths"). + Name(membershipStack.ClusterName). + Body(userPatch). + Do(context.Background()). + Error()).To(Succeed()) + + membershipClient.Orders() <- &generated.Order{ + Message: &generated.Order_ExistingStack{ + ExistingStack: membershipStack, + }, + } + + Eventually(func(g Gomega) { + a := &unstructured.Unstructured{} + g.Expect(LoadResource("Auths", membershipStack.ClusterName, a)).To(Succeed()) + g.Expect(a.GetAnnotations()).To(HaveKeyWithValue("user-custom/note", "do not delete")) + }).Should(Succeed()) + }) + }) When("removing modules", func() { var ( modulesToRemove map[string]struct{} From 84af91595417c20d11abad05001a7f9ec5ce1562 Mon Sep 17 00:00:00 2001 From: Geoffrey Ragot Date: Tue, 2 Jun 2026 09:31:26 +0200 Subject: [PATCH 2/3] test: add error path tests for critical sync functions Add unit tests using a mock K8SClient for: - syncExistingStack: Apply failure stops module sync, versions defaulting - syncModules: partial Apply failure continues to next module, delete failure - syncStargate: Apply/delete failures, nil StargateConfig - syncAuthClients: Apply failure skips client, List failure, stale delete failure - deleteStack: non-NotFound K8S error, NotFound sends message - disableStack/enableStack: Apply failure, field manager and force verification - createOrUpdate: conflict fallback to Get, conflict+Get failure, non-conflict error - informer_versions: Add/Update/Delete handlers, unchanged spec skipped - extractVersionsSpec: normal, nil, non-string values - generateMetadata: with/without labels and annotations Co-Authored-By: Claude Opus 4.6 (1M context) --- internal/error_paths_test.go | 861 +++++++++++++++++++++++++++++++ internal/k8s_client_mock_test.go | 63 +++ 2 files changed, 924 insertions(+) create mode 100644 internal/error_paths_test.go create mode 100644 internal/k8s_client_mock_test.go diff --git a/internal/error_paths_test.go b/internal/error_paths_test.go new file mode 100644 index 0000000..87b28e8 --- /dev/null +++ b/internal/error_paths_test.go @@ -0,0 +1,861 @@ +package internal + +import ( + "context" + "fmt" + "net/url" + "testing" + + v1apis "k8s.io/apiextensions-apiserver/pkg/apis/apiextensions/v1" + + "github.com/formancehq/go-libs/v2/logging" + "github.com/formancehq/stack/components/agent/internal/generated" + "github.com/google/uuid" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + apierrors "k8s.io/apimachinery/pkg/api/errors" + "k8s.io/apimachinery/pkg/api/meta" + "k8s.io/apimachinery/pkg/apis/meta/v1/unstructured" + "k8s.io/apimachinery/pkg/labels" + "k8s.io/apimachinery/pkg/runtime/schema" +) + +func newTestListener(k8s K8SClient, membership MembershipClient, modules []v1apis.CustomResourceDefinition) *membershipListener { + return NewMembershipListener( + k8s, + ClientInfo{ + BaseUrl: &url.URL{Scheme: "https", Host: "example.com"}, + }, + nil, // restMapper not needed when we mock Apply directly + membership, + modules, + ) +} + +func newTestListenerWithMapper(k8s K8SClient, membership MembershipClient, modules []v1apis.CustomResourceDefinition, mapper *fakeRESTMapper) *membershipListener { + return NewMembershipListener( + k8s, + ClientInfo{ + BaseUrl: &url.URL{Scheme: "https", Host: "example.com"}, + }, + mapper, + membership, + modules, + ) +} + +// fakeRESTMapper satisfies meta.RESTMapper for unit tests. +type fakeRESTMapper struct { + mappings map[schema.GroupKind]*fakeRESTMapping +} + +type fakeRESTMapping struct { + resource string +} + +func (f *fakeRESTMapper) KindFor(_ schema.GroupVersionResource) (schema.GroupVersionKind, error) { + return schema.GroupVersionKind{}, fmt.Errorf("not implemented") +} +func (f *fakeRESTMapper) KindsFor(_ schema.GroupVersionResource) ([]schema.GroupVersionKind, error) { + return nil, fmt.Errorf("not implemented") +} +func (f *fakeRESTMapper) ResourceFor(_ schema.GroupVersionResource) (schema.GroupVersionResource, error) { + return schema.GroupVersionResource{}, fmt.Errorf("not implemented") +} +func (f *fakeRESTMapper) ResourcesFor(_ schema.GroupVersionResource) ([]schema.GroupVersionResource, error) { + return nil, fmt.Errorf("not implemented") +} +func (f *fakeRESTMapper) RESTMapping(gk schema.GroupKind, _ ...string) (*meta.RESTMapping, error) { + m, ok := f.mappings[gk] + if !ok { + return nil, fmt.Errorf("no mapping for %s", gk) + } + return &meta.RESTMapping{ + Resource: schema.GroupVersionResource{ + Group: gk.Group, + Version: "v1beta1", + Resource: m.resource, + }, + GroupVersionKind: schema.GroupVersionKind{ + Group: gk.Group, + Version: "v1beta1", + Kind: gk.Kind, + }, + }, nil +} +func (f *fakeRESTMapper) RESTMappings(_ schema.GroupKind, _ ...string) ([]*meta.RESTMapping, error) { + return nil, fmt.Errorf("not implemented") +} +func (f *fakeRESTMapper) ResourceSingularizer(_ string) (string, error) { + return "", fmt.Errorf("not implemented") +} + +func defaultMapper() *fakeRESTMapper { + return &fakeRESTMapper{ + mappings: map[schema.GroupKind]*fakeRESTMapping{ + {Group: "formance.com", Kind: "Stack"}: {resource: "stacks"}, + {Group: "formance.com", Kind: "Auth"}: {resource: "auths"}, + {Group: "formance.com", Kind: "Gateway"}: {resource: "gateways"}, + {Group: "formance.com", Kind: "Ledger"}: {resource: "ledgers"}, + {Group: "formance.com", Kind: "Stargate"}: {resource: "stargates"}, + {Group: "formance.com", Kind: "AuthClient"}: {resource: "authclients"}, + {Group: "formance.com", Kind: "Webhooks"}: {resource: "webhooks"}, + {Group: "formance.com", Kind: "Payments"}: {resource: "payments"}, + {Group: "formance.com", Kind: "Search"}: {resource: "searches"}, + {Group: "formance.com", Kind: "Orchestration"}: {resource: "orchestrations"}, + {Group: "formance.com", Kind: "Wallets"}: {resource: "wallets"}, + {Group: "formance.com", Kind: "Reconciliation"}: {resource: "reconciliations"}, + }, + } +} + +func newFakeStack(name string) *unstructured.Unstructured { + s := &unstructured.Unstructured{} + s.SetName(name) + s.SetUID("fake-uid") + return s +} + +// --- syncExistingStack error paths --- + +func TestSyncExistingStack_ApplyFailure(t *testing.T) { + t.Parallel() + ctx := logging.TestingContext() + + applyErr := fmt.Errorf("connection refused") + k8s := &mockK8SClient{ + applyFn: func(_ context.Context, _ string, _ *unstructured.Unstructured, _ string, _ bool) (*unstructured.Unstructured, error) { + return nil, applyErr + }, + } + listener := newTestListenerWithMapper(k8s, NewMembershipClientMock(), nil, defaultMapper()) + + // Should not panic; error is logged, sync stops before modules + listener.syncExistingStack(ctx, &generated.Stack{ + ClusterName: uuid.NewString(), + AuthConfig: &generated.AuthConfig{}, + }) +} + +func TestSyncExistingStack_VersionsDefault(t *testing.T) { + t.Parallel() + ctx := logging.TestingContext() + + var appliedContent map[string]any + k8s := &mockK8SClient{ + applyFn: func(_ context.Context, resource string, obj *unstructured.Unstructured, _ string, _ bool) (*unstructured.Unstructured, error) { + if resource == "stacks" { + appliedContent = obj.Object + } + result := obj.DeepCopy() + result.SetUID("fake-uid") + return result, nil + }, + } + listener := newTestListenerWithMapper(k8s, NewMembershipClientMock(), nil, defaultMapper()) + + listener.syncExistingStack(ctx, &generated.Stack{ + ClusterName: uuid.NewString(), + Versions: "", + AuthConfig: &generated.AuthConfig{}, + }) + + spec, _ := appliedContent["spec"].(map[string]any) + require.Equal(t, "default", spec["versionsFromFile"]) +} + +func TestSyncExistingStack_ExplicitVersions(t *testing.T) { + t.Parallel() + ctx := logging.TestingContext() + + var appliedContent map[string]any + k8s := &mockK8SClient{ + applyFn: func(_ context.Context, resource string, obj *unstructured.Unstructured, _ string, _ bool) (*unstructured.Unstructured, error) { + if resource == "stacks" { + appliedContent = obj.Object + } + result := obj.DeepCopy() + result.SetUID("fake-uid") + return result, nil + }, + } + listener := newTestListenerWithMapper(k8s, NewMembershipClientMock(), nil, defaultMapper()) + + listener.syncExistingStack(ctx, &generated.Stack{ + ClusterName: uuid.NewString(), + Versions: "v1.2.3", + AuthConfig: &generated.AuthConfig{}, + }) + + spec, _ := appliedContent["spec"].(map[string]any) + require.Equal(t, "v1.2.3", spec["versionsFromFile"]) +} + +// --- syncModules error paths --- + +func testModuleCRDs() []v1apis.CustomResourceDefinition { + makeCRD := func(kind, singular, plural string) v1apis.CustomResourceDefinition { + return v1apis.CustomResourceDefinition{ + Spec: v1apis.CustomResourceDefinitionSpec{ + Group: "formance.com", + Names: v1apis.CustomResourceDefinitionNames{Kind: kind}, + Versions: []v1apis.CustomResourceDefinitionVersion{ + {Name: "v1beta1"}, + }, + }, + Status: v1apis.CustomResourceDefinitionStatus{ + AcceptedNames: v1apis.CustomResourceDefinitionNames{ + Singular: singular, + Plural: plural, + }, + }, + } + } + return []v1apis.CustomResourceDefinition{ + makeCRD("Auth", "auth", "auths"), + makeCRD("Gateway", "gateway", "gateways"), + makeCRD("Ledger", "ledger", "ledgers"), + } +} + +func TestSyncModules_PartialApplyFailure(t *testing.T) { + t.Parallel() + ctx := logging.TestingContext() + + var appliedResources []string + k8s := &mockK8SClient{ + applyFn: func(_ context.Context, resource string, obj *unstructured.Unstructured, _ string, _ bool) (*unstructured.Unstructured, error) { + appliedResources = append(appliedResources, resource) + if resource == "auths" { + return nil, fmt.Errorf("auth apply failed") + } + result := obj.DeepCopy() + result.SetUID("fake-uid") + return result, nil + }, + } + + modules := testModuleCRDs() + listener := newTestListenerWithMapper(k8s, NewMembershipClientMock(), modules, defaultMapper()) + + stack := newFakeStack("test-stack") + membershipStack := &generated.Stack{ + ClusterName: "test-stack", + AuthConfig: &generated.AuthConfig{}, + Modules: []*generated.Module{ + {Name: "Auth"}, + {Name: "Gateway"}, + {Name: "Ledger"}, + }, + } + + // Should not panic; Auth fails but Gateway and Ledger are still attempted + listener.syncModules(ctx, map[string]any{}, stack, membershipStack) + + assert.Contains(t, appliedResources, "auths") + assert.Contains(t, appliedResources, "gateways") + assert.Contains(t, appliedResources, "ledgers") +} + +func TestSyncModules_DeleteModuleFailure(t *testing.T) { + t.Parallel() + ctx := logging.TestingContext() + + deleteErr := fmt.Errorf("delete forbidden") + k8s := &mockK8SClient{ + ensureNotExistsBySelectorFn: func(_ context.Context, _ string, _ labels.Selector) error { + return deleteErr + }, + } + + modules := testModuleCRDs() + listener := newTestListenerWithMapper(k8s, NewMembershipClientMock(), modules, defaultMapper()) + + stack := newFakeStack("test-stack") + membershipStack := &generated.Stack{ + ClusterName: "test-stack", + AuthConfig: &generated.AuthConfig{}, + Modules: []*generated.Module{}, // no modules expected → all should be deleted + } + + // Should not panic; errors logged, loop continues + listener.syncModules(ctx, map[string]any{}, stack, membershipStack) +} + +// --- syncStargate error paths --- + +func TestSyncStargate_ApplyFailure(t *testing.T) { + t.Parallel() + ctx := logging.TestingContext() + + k8s := &mockK8SClient{ + applyFn: func(_ context.Context, _ string, _ *unstructured.Unstructured, _ string, _ bool) (*unstructured.Unstructured, error) { + return nil, fmt.Errorf("stargate apply failed") + }, + } + listener := newTestListenerWithMapper(k8s, NewMembershipClientMock(), nil, defaultMapper()) + + stack := newFakeStack("org-stack") + membershipStack := &generated.Stack{ + AuthConfig: &generated.AuthConfig{}, + StargateConfig: &generated.StargateConfig{Enabled: true}, + } + + // Should not panic + listener.syncStargate(ctx, map[string]any{}, stack, membershipStack) +} + +func TestSyncStargate_DeleteFailure(t *testing.T) { + t.Parallel() + ctx := logging.TestingContext() + + k8s := &mockK8SClient{ + ensureNotExistsFn: func(_ context.Context, _, _ string) error { + return fmt.Errorf("delete failed") + }, + } + listener := newTestListenerWithMapper(k8s, NewMembershipClientMock(), nil, defaultMapper()) + + stack := newFakeStack("org-stack") + membershipStack := &generated.Stack{ + StargateConfig: &generated.StargateConfig{Enabled: false}, + } + + // Should not panic + listener.syncStargate(ctx, map[string]any{}, stack, membershipStack) +} + +func TestSyncStargate_NilConfig(t *testing.T) { + t.Parallel() + ctx := logging.TestingContext() + + var deleteCalled bool + k8s := &mockK8SClient{ + ensureNotExistsFn: func(_ context.Context, _, _ string) error { + deleteCalled = true + return nil + }, + } + listener := newTestListenerWithMapper(k8s, NewMembershipClientMock(), nil, defaultMapper()) + + stack := newFakeStack("org-stack") + membershipStack := &generated.Stack{ + StargateConfig: nil, + } + + listener.syncStargate(ctx, map[string]any{}, stack, membershipStack) + assert.True(t, deleteCalled, "should attempt to delete stargate when config is nil") +} + +// --- syncAuthClients error paths --- + +func TestSyncAuthClients_ApplyFailure(t *testing.T) { + t.Parallel() + ctx := logging.TestingContext() + + callCount := 0 + k8s := &mockK8SClient{ + applyFn: func(_ context.Context, _ string, obj *unstructured.Unstructured, _ string, _ bool) (*unstructured.Unstructured, error) { + callCount++ + if callCount == 1 { + return nil, fmt.Errorf("apply failed for first client") + } + result := obj.DeepCopy() + result.SetUID("fake-uid") + return result, nil + }, + } + listener := newTestListenerWithMapper(k8s, NewMembershipClientMock(), nil, defaultMapper()) + + stack := newFakeStack("test-stack") + clients := []*generated.AuthClient{ + {Id: "client1", Public: true}, + {Id: "client2", Public: false}, + } + + // Should not panic; first client fails, second succeeds + listener.syncAuthClients(ctx, map[string]any{}, stack, clients) + assert.Equal(t, 2, callCount, "should attempt apply for both clients") +} + +func TestSyncAuthClients_ListFailure(t *testing.T) { + t.Parallel() + ctx := logging.TestingContext() + + k8s := &mockK8SClient{ + applyFn: func(_ context.Context, _ string, obj *unstructured.Unstructured, _ string, _ bool) (*unstructured.Unstructured, error) { + result := obj.DeepCopy() + result.SetUID("fake-uid") + return result, nil + }, + listFn: func(_ context.Context, _ string, _ labels.Selector) ([]unstructured.Unstructured, error) { + return nil, fmt.Errorf("list forbidden") + }, + } + listener := newTestListenerWithMapper(k8s, NewMembershipClientMock(), nil, defaultMapper()) + + stack := newFakeStack("test-stack") + clients := []*generated.AuthClient{ + {Id: "client1", Public: true}, + } + + // Should not panic; list fails so stale clients are not cleaned up + listener.syncAuthClients(ctx, map[string]any{}, stack, clients) +} + +func TestSyncAuthClients_DeleteStaleFailure(t *testing.T) { + t.Parallel() + ctx := logging.TestingContext() + + var deleteAttempted []string + k8s := &mockK8SClient{ + applyFn: func(_ context.Context, _ string, obj *unstructured.Unstructured, _ string, _ bool) (*unstructured.Unstructured, error) { + result := obj.DeepCopy() + result.SetUID("fake-uid") + return result, nil + }, + listFn: func(_ context.Context, _ string, _ labels.Selector) ([]unstructured.Unstructured, error) { + return []unstructured.Unstructured{ + {Object: map[string]any{"metadata": map[string]any{"name": "test-stack-client1"}}}, + {Object: map[string]any{"metadata": map[string]any{"name": "test-stack-stale"}}}, + }, nil + }, + ensureNotExistsFn: func(_ context.Context, _ string, name string) error { + deleteAttempted = append(deleteAttempted, name) + return fmt.Errorf("delete forbidden") + }, + } + listener := newTestListenerWithMapper(k8s, NewMembershipClientMock(), nil, defaultMapper()) + + stack := newFakeStack("test-stack") + clients := []*generated.AuthClient{ + {Id: "client1", Public: true}, + } + + // Should not panic; stale client deletion fails but is logged + listener.syncAuthClients(ctx, map[string]any{}, stack, clients) + assert.Contains(t, deleteAttempted, "test-stack-stale") +} + +// --- deleteStack error paths --- + +func TestDeleteStack_K8SError(t *testing.T) { + t.Parallel() + ctx := logging.TestingContext() + + k8s := &mockK8SClient{ + deleteFn: func(_ context.Context, _, _ string) error { + return fmt.Errorf("internal server error") + }, + } + mock := NewMembershipClientMock() + listener := newTestListener(k8s, mock, nil) + + // Should not panic; error logged, no message sent to membership + listener.deleteStack(ctx, &generated.DeletedStack{ClusterName: "broken-stack"}) + assert.Empty(t, mock.GetMessages(), "should not send any message on non-NotFound error") +} + +func TestDeleteStack_SendFailure(t *testing.T) { + t.Parallel() + ctx := logging.TestingContext() + + k8s := &mockK8SClient{ + deleteFn: func(_ context.Context, _, _ string) error { + return apierrors.NewNotFound(schema.GroupResource{Group: "formance.com", Resource: "stacks"}, "gone-stack") + }, + } + + // Use a real mock that will track messages, but the test verifies no panic + mock := NewMembershipClientMock() + listener := newTestListener(k8s, mock, nil) + + listener.deleteStack(ctx, &generated.DeletedStack{ClusterName: "gone-stack"}) + require.Len(t, mock.GetMessages(), 1) + assert.NotNil(t, mock.GetMessages()[0].GetStackDeleted()) +} + +// --- disableStack / enableStack error paths --- + +func TestDisableStack_ApplyFailure(t *testing.T) { + t.Parallel() + ctx := logging.TestingContext() + + k8s := &mockK8SClient{ + applyFn: func(_ context.Context, _ string, _ *unstructured.Unstructured, _ string, _ bool) (*unstructured.Unstructured, error) { + return nil, fmt.Errorf("apply failed") + }, + } + listener := newTestListener(k8s, NewMembershipClientMock(), nil) + + // Should not panic + listener.disableStack(ctx, &generated.DisabledStack{ClusterName: "test-stack"}) +} + +func TestEnableStack_ApplyFailure(t *testing.T) { + t.Parallel() + ctx := logging.TestingContext() + + k8s := &mockK8SClient{ + applyFn: func(_ context.Context, _ string, _ *unstructured.Unstructured, _ string, _ bool) (*unstructured.Unstructured, error) { + return nil, fmt.Errorf("apply failed") + }, + } + listener := newTestListener(k8s, NewMembershipClientMock(), nil) + + // Should not panic + listener.enableStack(ctx, &generated.EnabledStack{ClusterName: "test-stack"}) +} + +func TestDisableStack_FieldManager(t *testing.T) { + t.Parallel() + ctx := logging.TestingContext() + + var capturedFieldManager string + var capturedForce bool + k8s := &mockK8SClient{ + applyFn: func(_ context.Context, _ string, _ *unstructured.Unstructured, fm string, force bool) (*unstructured.Unstructured, error) { + capturedFieldManager = fm + capturedForce = force + return &unstructured.Unstructured{}, nil + }, + } + listener := newTestListener(k8s, NewMembershipClientMock(), nil) + + listener.disableStack(ctx, &generated.DisabledStack{ClusterName: "test-stack"}) + assert.Equal(t, fieldManagerDisable, capturedFieldManager) + assert.True(t, capturedForce, "disable should use force=true") +} + +func TestEnableStack_FieldManager(t *testing.T) { + t.Parallel() + ctx := logging.TestingContext() + + var capturedFieldManager string + var capturedForce bool + k8s := &mockK8SClient{ + applyFn: func(_ context.Context, _ string, _ *unstructured.Unstructured, fm string, force bool) (*unstructured.Unstructured, error) { + capturedFieldManager = fm + capturedForce = force + return &unstructured.Unstructured{}, nil + }, + } + listener := newTestListener(k8s, NewMembershipClientMock(), nil) + + listener.enableStack(ctx, &generated.EnabledStack{ClusterName: "test-stack"}) + assert.Equal(t, fieldManagerDisable, capturedFieldManager) + assert.True(t, capturedForce, "enable should use force=true") +} + +// --- createOrUpdate error paths --- + +func TestCreateOrUpdate_ConflictFallsBackToGet(t *testing.T) { + t.Parallel() + ctx := logging.TestingContext() + + existingObj := &unstructured.Unstructured{ + Object: map[string]any{ + "metadata": map[string]any{ + "name": "my-stack", + "annotations": map[string]any{ + "user/custom": "preserved", + }, + }, + }, + } + + k8s := &mockK8SClient{ + applyFn: func(_ context.Context, _ string, _ *unstructured.Unstructured, _ string, _ bool) (*unstructured.Unstructured, error) { + return nil, apierrors.NewConflict(schema.GroupResource{Group: "formance.com", Resource: "stacks"}, "my-stack", fmt.Errorf("field owned by kubectl-edit")) + }, + getFn: func(_ context.Context, _ string, _ string) (*unstructured.Unstructured, error) { + return existingObj, nil + }, + } + + listener := newTestListenerWithMapper(k8s, NewMembershipClientMock(), nil, defaultMapper()) + + result, err := listener.createOrUpdate(ctx, formanceGroupVersion.WithKind("Stack"), "my-stack", "my-stack", nil, map[string]any{ + "spec": map[string]any{"versionsFromFile": "default"}, + }) + + require.NoError(t, err) + assert.Equal(t, "preserved", result.GetAnnotations()["user/custom"]) +} + +func TestCreateOrUpdate_ConflictGetFailure(t *testing.T) { + t.Parallel() + ctx := logging.TestingContext() + + k8s := &mockK8SClient{ + applyFn: func(_ context.Context, _ string, _ *unstructured.Unstructured, _ string, _ bool) (*unstructured.Unstructured, error) { + return nil, apierrors.NewConflict(schema.GroupResource{}, "x", fmt.Errorf("conflict")) + }, + getFn: func(_ context.Context, _ string, _ string) (*unstructured.Unstructured, error) { + return nil, fmt.Errorf("get also failed") + }, + } + + listener := newTestListenerWithMapper(k8s, NewMembershipClientMock(), nil, defaultMapper()) + + _, err := listener.createOrUpdate(ctx, formanceGroupVersion.WithKind("Stack"), "x", "x", nil, map[string]any{ + "spec": map[string]any{}, + }) + + require.Error(t, err) + assert.Contains(t, err.Error(), "getting existing object after conflict") +} + +func TestCreateOrUpdate_NonConflictError(t *testing.T) { + t.Parallel() + ctx := logging.TestingContext() + + k8s := &mockK8SClient{ + applyFn: func(_ context.Context, _ string, _ *unstructured.Unstructured, _ string, _ bool) (*unstructured.Unstructured, error) { + return nil, fmt.Errorf("network error") + }, + } + + listener := newTestListenerWithMapper(k8s, NewMembershipClientMock(), nil, defaultMapper()) + + _, err := listener.createOrUpdate(ctx, formanceGroupVersion.WithKind("Stack"), "x", "x", nil, map[string]any{ + "spec": map[string]any{}, + }) + + require.Error(t, err) + assert.Contains(t, err.Error(), "applying object") +} + +func TestCreateOrUpdate_FieldManagerIsAgent(t *testing.T) { + t.Parallel() + ctx := logging.TestingContext() + + var capturedFM string + var capturedForce bool + k8s := &mockK8SClient{ + applyFn: func(_ context.Context, _ string, obj *unstructured.Unstructured, fm string, force bool) (*unstructured.Unstructured, error) { + capturedFM = fm + capturedForce = force + result := obj.DeepCopy() + result.SetUID("fake-uid") + return result, nil + }, + } + + listener := newTestListenerWithMapper(k8s, NewMembershipClientMock(), nil, defaultMapper()) + + _, err := listener.createOrUpdate(ctx, formanceGroupVersion.WithKind("Stack"), "x", "x", nil, map[string]any{ + "spec": map[string]any{}, + }) + + require.NoError(t, err) + assert.Equal(t, fieldManagerAgent, capturedFM) + assert.False(t, capturedForce, "createOrUpdate should use force=false") +} + +// --- informer_versions (entirely untested until now) --- + +func TestVersionsEventHandler_Add(t *testing.T) { + t.Parallel() + + mock := NewMembershipClientMock() + handler := VersionsEventHandler(logging.Testing(), mock) + + version := &unstructured.Unstructured{ + Object: map[string]any{ + "metadata": map[string]any{ + "name": "v2.0", + "annotations": map[string]any{ + "formance.com/deprecated": "true", + }, + }, + "spec": map[string]any{ + "ledger": "v2.0.0", + "payments": "v1.5.0", + }, + }, + } + + handler.OnAdd(version, false) + + require.Len(t, mock.GetMessages(), 1) + msg := mock.GetMessages()[0].GetAddedVersion() + require.NotNil(t, msg) + assert.Equal(t, "v2.0", msg.Name) + assert.Equal(t, "v2.0.0", msg.Versions["ledger"]) + assert.Equal(t, "v1.5.0", msg.Versions["payments"]) + assert.True(t, msg.Deprecated) +} + +func TestVersionsEventHandler_Add_NoSpec(t *testing.T) { + t.Parallel() + + mock := NewMembershipClientMock() + handler := VersionsEventHandler(logging.Testing(), mock) + + version := &unstructured.Unstructured{ + Object: map[string]any{ + "metadata": map[string]any{ + "name": "empty", + }, + }, + } + + handler.OnAdd(version, false) + + require.Len(t, mock.GetMessages(), 1) + msg := mock.GetMessages()[0].GetAddedVersion() + require.NotNil(t, msg) + assert.Nil(t, msg.Versions) + assert.False(t, msg.Deprecated) +} + +func TestVersionsEventHandler_Update_Changed(t *testing.T) { + t.Parallel() + + mock := NewMembershipClientMock() + handler := VersionsEventHandler(logging.Testing(), mock) + + oldVersion := &unstructured.Unstructured{ + Object: map[string]any{ + "metadata": map[string]any{"name": "v1"}, + "spec": map[string]any{"ledger": "v1.0.0"}, + }, + } + newVersion := &unstructured.Unstructured{ + Object: map[string]any{ + "metadata": map[string]any{"name": "v1"}, + "spec": map[string]any{"ledger": "v1.1.0"}, + }, + } + + handler.OnUpdate(oldVersion, newVersion) + + require.Len(t, mock.GetMessages(), 1) + msg := mock.GetMessages()[0].GetUpdatedVersion() + require.NotNil(t, msg) + assert.Equal(t, "v1.1.0", msg.Versions["ledger"]) +} + +func TestVersionsEventHandler_Update_Unchanged(t *testing.T) { + t.Parallel() + + mock := NewMembershipClientMock() + handler := VersionsEventHandler(logging.Testing(), mock) + + obj := &unstructured.Unstructured{ + Object: map[string]any{ + "metadata": map[string]any{"name": "v1"}, + "spec": map[string]any{"ledger": "v1.0.0"}, + }, + } + + handler.OnUpdate(obj, obj) + + assert.Empty(t, mock.GetMessages(), "should not send message when spec unchanged") +} + +func TestVersionsEventHandler_Delete(t *testing.T) { + t.Parallel() + + mock := NewMembershipClientMock() + handler := VersionsEventHandler(logging.Testing(), mock) + + version := &unstructured.Unstructured{ + Object: map[string]any{ + "metadata": map[string]any{"name": "v1"}, + }, + } + + handler.OnDelete(version) + + require.Len(t, mock.GetMessages(), 1) + msg := mock.GetMessages()[0].GetDeletedVersion() + require.NotNil(t, msg) + assert.Equal(t, "v1", msg.Name) +} + +// --- extractVersionsSpec --- + +func TestExtractVersionsSpec(t *testing.T) { + t.Parallel() + + tests := []struct { + name string + obj *unstructured.Unstructured + expect map[string]string + }{ + { + name: "normal spec", + obj: &unstructured.Unstructured{ + Object: map[string]any{ + "spec": map[string]any{ + "ledger": "v1.0", "payments": "v2.0", + }, + }, + }, + expect: map[string]string{"ledger": "v1.0", "payments": "v2.0"}, + }, + { + name: "nil spec", + obj: &unstructured.Unstructured{ + Object: map[string]any{}, + }, + expect: nil, + }, + { + name: "non-string values ignored", + obj: &unstructured.Unstructured{ + Object: map[string]any{ + "spec": map[string]any{ + "ledger": "v1.0", + "count": int64(42), + }, + }, + }, + expect: map[string]string{"ledger": "v1.0"}, + }, + } + + for _, tc := range tests { + t.Run(tc.name, func(t *testing.T) { + t.Parallel() + result := extractVersionsSpec(tc.obj) + assert.Equal(t, tc.expect, result) + }) + } +} + +// --- generateMetadata --- + +func TestGenerateMetadata(t *testing.T) { + t.Parallel() + + listener := newTestListener(&mockK8SClient{}, NewMembershipClientMock(), nil) + + t.Run("with labels and annotations", func(t *testing.T) { + stack := &generated.Stack{ + AdditionalLabels: map[string]string{"env": "prod", "team": "infra"}, + AdditionalAnnotations: map[string]string{"note": "important"}, + } + md := listener.generateMetadata(stack) + + labels := md["labels"].(map[string]any) + assert.Equal(t, "prod", labels["formance.com/env"]) + assert.Equal(t, "infra", labels["formance.com/team"]) + + annotations := md["annotations"].(map[string]any) + assert.Equal(t, "important", annotations["formance.com/note"]) + }) + + t.Run("nil labels and annotations", func(t *testing.T) { + stack := &generated.Stack{} + md := listener.generateMetadata(stack) + + labels := md["labels"].(map[string]any) + assert.Empty(t, labels) + + annotations := md["annotations"].(map[string]any) + assert.Empty(t, annotations) + }) +} diff --git a/internal/k8s_client_mock_test.go b/internal/k8s_client_mock_test.go new file mode 100644 index 0000000..5d61216 --- /dev/null +++ b/internal/k8s_client_mock_test.go @@ -0,0 +1,63 @@ +package internal + +import ( + "context" + + "k8s.io/apimachinery/pkg/apis/meta/v1/unstructured" + "k8s.io/apimachinery/pkg/labels" +) + +// mockK8SClient is a configurable mock for unit-testing error paths. +// Each field is a function that, when set, overrides the default (no-op) behaviour. +type mockK8SClient struct { + getFn func(ctx context.Context, resource, name string) (*unstructured.Unstructured, error) + applyFn func(ctx context.Context, resource string, obj *unstructured.Unstructured, fieldManager string, force bool) (*unstructured.Unstructured, error) + deleteFn func(ctx context.Context, resource, name string) error + ensureNotExistsFn func(ctx context.Context, resource, name string) error + ensureNotExistsBySelectorFn func(ctx context.Context, resource string, selector labels.Selector) error + listFn func(ctx context.Context, resource string, selector labels.Selector) ([]unstructured.Unstructured, error) +} + +func (m *mockK8SClient) Get(ctx context.Context, resource, name string) (*unstructured.Unstructured, error) { + if m.getFn != nil { + return m.getFn(ctx, resource, name) + } + return &unstructured.Unstructured{}, nil +} + +func (m *mockK8SClient) Apply(ctx context.Context, resource string, obj *unstructured.Unstructured, fieldManager string, force bool) (*unstructured.Unstructured, error) { + if m.applyFn != nil { + return m.applyFn(ctx, resource, obj, fieldManager, force) + } + result := obj.DeepCopy() + result.SetUID("fake-uid") + return result, nil +} + +func (m *mockK8SClient) Delete(ctx context.Context, resource, name string) error { + if m.deleteFn != nil { + return m.deleteFn(ctx, resource, name) + } + return nil +} + +func (m *mockK8SClient) EnsureNotExists(ctx context.Context, resource, name string) error { + if m.ensureNotExistsFn != nil { + return m.ensureNotExistsFn(ctx, resource, name) + } + return nil +} + +func (m *mockK8SClient) EnsureNotExistsBySelector(ctx context.Context, resource string, selector labels.Selector) error { + if m.ensureNotExistsBySelectorFn != nil { + return m.ensureNotExistsBySelectorFn(ctx, resource, selector) + } + return nil +} + +func (m *mockK8SClient) List(ctx context.Context, resource string, selector labels.Selector) ([]unstructured.Unstructured, error) { + if m.listFn != nil { + return m.listFn(ctx, resource, selector) + } + return nil, nil +} From 203eec76c2c9503e850e75d6997a0cc01b80df86 Mon Sep 17 00:00:00 2001 From: Geoffrey Ragot Date: Tue, 2 Jun 2026 09:58:16 +0200 Subject: [PATCH 3/3] test: replace mock-based tests with integration tests Remove mock-based unit tests and add proper integration tests using envtest, matching the project's testing conventions: Membership listener: - Verify versionsFromFile defaults to "default" when empty - Verify explicit versions are propagated - Verify SSA field manager "formance-agent" appears in managedFields - Verify resync is idempotent (same resourceVersion after re-apply) Versions informer (previously entirely untested): - Add/Update/Delete handlers with real K8S resources - Spec extraction with deprecated annotation - Empty spec handling - Update only sent when spec actually changes Co-Authored-By: Claude Opus 4.6 (1M context) --- internal/error_paths_test.go | 861 ------------------------------ internal/k8s_client_mock_test.go | 63 --- tests/membership_listener_test.go | 56 ++ tests/versions_informer_test.go | 192 +++++++ 4 files changed, 248 insertions(+), 924 deletions(-) delete mode 100644 internal/error_paths_test.go delete mode 100644 internal/k8s_client_mock_test.go create mode 100644 tests/versions_informer_test.go diff --git a/internal/error_paths_test.go b/internal/error_paths_test.go deleted file mode 100644 index 87b28e8..0000000 --- a/internal/error_paths_test.go +++ /dev/null @@ -1,861 +0,0 @@ -package internal - -import ( - "context" - "fmt" - "net/url" - "testing" - - v1apis "k8s.io/apiextensions-apiserver/pkg/apis/apiextensions/v1" - - "github.com/formancehq/go-libs/v2/logging" - "github.com/formancehq/stack/components/agent/internal/generated" - "github.com/google/uuid" - "github.com/stretchr/testify/assert" - "github.com/stretchr/testify/require" - apierrors "k8s.io/apimachinery/pkg/api/errors" - "k8s.io/apimachinery/pkg/api/meta" - "k8s.io/apimachinery/pkg/apis/meta/v1/unstructured" - "k8s.io/apimachinery/pkg/labels" - "k8s.io/apimachinery/pkg/runtime/schema" -) - -func newTestListener(k8s K8SClient, membership MembershipClient, modules []v1apis.CustomResourceDefinition) *membershipListener { - return NewMembershipListener( - k8s, - ClientInfo{ - BaseUrl: &url.URL{Scheme: "https", Host: "example.com"}, - }, - nil, // restMapper not needed when we mock Apply directly - membership, - modules, - ) -} - -func newTestListenerWithMapper(k8s K8SClient, membership MembershipClient, modules []v1apis.CustomResourceDefinition, mapper *fakeRESTMapper) *membershipListener { - return NewMembershipListener( - k8s, - ClientInfo{ - BaseUrl: &url.URL{Scheme: "https", Host: "example.com"}, - }, - mapper, - membership, - modules, - ) -} - -// fakeRESTMapper satisfies meta.RESTMapper for unit tests. -type fakeRESTMapper struct { - mappings map[schema.GroupKind]*fakeRESTMapping -} - -type fakeRESTMapping struct { - resource string -} - -func (f *fakeRESTMapper) KindFor(_ schema.GroupVersionResource) (schema.GroupVersionKind, error) { - return schema.GroupVersionKind{}, fmt.Errorf("not implemented") -} -func (f *fakeRESTMapper) KindsFor(_ schema.GroupVersionResource) ([]schema.GroupVersionKind, error) { - return nil, fmt.Errorf("not implemented") -} -func (f *fakeRESTMapper) ResourceFor(_ schema.GroupVersionResource) (schema.GroupVersionResource, error) { - return schema.GroupVersionResource{}, fmt.Errorf("not implemented") -} -func (f *fakeRESTMapper) ResourcesFor(_ schema.GroupVersionResource) ([]schema.GroupVersionResource, error) { - return nil, fmt.Errorf("not implemented") -} -func (f *fakeRESTMapper) RESTMapping(gk schema.GroupKind, _ ...string) (*meta.RESTMapping, error) { - m, ok := f.mappings[gk] - if !ok { - return nil, fmt.Errorf("no mapping for %s", gk) - } - return &meta.RESTMapping{ - Resource: schema.GroupVersionResource{ - Group: gk.Group, - Version: "v1beta1", - Resource: m.resource, - }, - GroupVersionKind: schema.GroupVersionKind{ - Group: gk.Group, - Version: "v1beta1", - Kind: gk.Kind, - }, - }, nil -} -func (f *fakeRESTMapper) RESTMappings(_ schema.GroupKind, _ ...string) ([]*meta.RESTMapping, error) { - return nil, fmt.Errorf("not implemented") -} -func (f *fakeRESTMapper) ResourceSingularizer(_ string) (string, error) { - return "", fmt.Errorf("not implemented") -} - -func defaultMapper() *fakeRESTMapper { - return &fakeRESTMapper{ - mappings: map[schema.GroupKind]*fakeRESTMapping{ - {Group: "formance.com", Kind: "Stack"}: {resource: "stacks"}, - {Group: "formance.com", Kind: "Auth"}: {resource: "auths"}, - {Group: "formance.com", Kind: "Gateway"}: {resource: "gateways"}, - {Group: "formance.com", Kind: "Ledger"}: {resource: "ledgers"}, - {Group: "formance.com", Kind: "Stargate"}: {resource: "stargates"}, - {Group: "formance.com", Kind: "AuthClient"}: {resource: "authclients"}, - {Group: "formance.com", Kind: "Webhooks"}: {resource: "webhooks"}, - {Group: "formance.com", Kind: "Payments"}: {resource: "payments"}, - {Group: "formance.com", Kind: "Search"}: {resource: "searches"}, - {Group: "formance.com", Kind: "Orchestration"}: {resource: "orchestrations"}, - {Group: "formance.com", Kind: "Wallets"}: {resource: "wallets"}, - {Group: "formance.com", Kind: "Reconciliation"}: {resource: "reconciliations"}, - }, - } -} - -func newFakeStack(name string) *unstructured.Unstructured { - s := &unstructured.Unstructured{} - s.SetName(name) - s.SetUID("fake-uid") - return s -} - -// --- syncExistingStack error paths --- - -func TestSyncExistingStack_ApplyFailure(t *testing.T) { - t.Parallel() - ctx := logging.TestingContext() - - applyErr := fmt.Errorf("connection refused") - k8s := &mockK8SClient{ - applyFn: func(_ context.Context, _ string, _ *unstructured.Unstructured, _ string, _ bool) (*unstructured.Unstructured, error) { - return nil, applyErr - }, - } - listener := newTestListenerWithMapper(k8s, NewMembershipClientMock(), nil, defaultMapper()) - - // Should not panic; error is logged, sync stops before modules - listener.syncExistingStack(ctx, &generated.Stack{ - ClusterName: uuid.NewString(), - AuthConfig: &generated.AuthConfig{}, - }) -} - -func TestSyncExistingStack_VersionsDefault(t *testing.T) { - t.Parallel() - ctx := logging.TestingContext() - - var appliedContent map[string]any - k8s := &mockK8SClient{ - applyFn: func(_ context.Context, resource string, obj *unstructured.Unstructured, _ string, _ bool) (*unstructured.Unstructured, error) { - if resource == "stacks" { - appliedContent = obj.Object - } - result := obj.DeepCopy() - result.SetUID("fake-uid") - return result, nil - }, - } - listener := newTestListenerWithMapper(k8s, NewMembershipClientMock(), nil, defaultMapper()) - - listener.syncExistingStack(ctx, &generated.Stack{ - ClusterName: uuid.NewString(), - Versions: "", - AuthConfig: &generated.AuthConfig{}, - }) - - spec, _ := appliedContent["spec"].(map[string]any) - require.Equal(t, "default", spec["versionsFromFile"]) -} - -func TestSyncExistingStack_ExplicitVersions(t *testing.T) { - t.Parallel() - ctx := logging.TestingContext() - - var appliedContent map[string]any - k8s := &mockK8SClient{ - applyFn: func(_ context.Context, resource string, obj *unstructured.Unstructured, _ string, _ bool) (*unstructured.Unstructured, error) { - if resource == "stacks" { - appliedContent = obj.Object - } - result := obj.DeepCopy() - result.SetUID("fake-uid") - return result, nil - }, - } - listener := newTestListenerWithMapper(k8s, NewMembershipClientMock(), nil, defaultMapper()) - - listener.syncExistingStack(ctx, &generated.Stack{ - ClusterName: uuid.NewString(), - Versions: "v1.2.3", - AuthConfig: &generated.AuthConfig{}, - }) - - spec, _ := appliedContent["spec"].(map[string]any) - require.Equal(t, "v1.2.3", spec["versionsFromFile"]) -} - -// --- syncModules error paths --- - -func testModuleCRDs() []v1apis.CustomResourceDefinition { - makeCRD := func(kind, singular, plural string) v1apis.CustomResourceDefinition { - return v1apis.CustomResourceDefinition{ - Spec: v1apis.CustomResourceDefinitionSpec{ - Group: "formance.com", - Names: v1apis.CustomResourceDefinitionNames{Kind: kind}, - Versions: []v1apis.CustomResourceDefinitionVersion{ - {Name: "v1beta1"}, - }, - }, - Status: v1apis.CustomResourceDefinitionStatus{ - AcceptedNames: v1apis.CustomResourceDefinitionNames{ - Singular: singular, - Plural: plural, - }, - }, - } - } - return []v1apis.CustomResourceDefinition{ - makeCRD("Auth", "auth", "auths"), - makeCRD("Gateway", "gateway", "gateways"), - makeCRD("Ledger", "ledger", "ledgers"), - } -} - -func TestSyncModules_PartialApplyFailure(t *testing.T) { - t.Parallel() - ctx := logging.TestingContext() - - var appliedResources []string - k8s := &mockK8SClient{ - applyFn: func(_ context.Context, resource string, obj *unstructured.Unstructured, _ string, _ bool) (*unstructured.Unstructured, error) { - appliedResources = append(appliedResources, resource) - if resource == "auths" { - return nil, fmt.Errorf("auth apply failed") - } - result := obj.DeepCopy() - result.SetUID("fake-uid") - return result, nil - }, - } - - modules := testModuleCRDs() - listener := newTestListenerWithMapper(k8s, NewMembershipClientMock(), modules, defaultMapper()) - - stack := newFakeStack("test-stack") - membershipStack := &generated.Stack{ - ClusterName: "test-stack", - AuthConfig: &generated.AuthConfig{}, - Modules: []*generated.Module{ - {Name: "Auth"}, - {Name: "Gateway"}, - {Name: "Ledger"}, - }, - } - - // Should not panic; Auth fails but Gateway and Ledger are still attempted - listener.syncModules(ctx, map[string]any{}, stack, membershipStack) - - assert.Contains(t, appliedResources, "auths") - assert.Contains(t, appliedResources, "gateways") - assert.Contains(t, appliedResources, "ledgers") -} - -func TestSyncModules_DeleteModuleFailure(t *testing.T) { - t.Parallel() - ctx := logging.TestingContext() - - deleteErr := fmt.Errorf("delete forbidden") - k8s := &mockK8SClient{ - ensureNotExistsBySelectorFn: func(_ context.Context, _ string, _ labels.Selector) error { - return deleteErr - }, - } - - modules := testModuleCRDs() - listener := newTestListenerWithMapper(k8s, NewMembershipClientMock(), modules, defaultMapper()) - - stack := newFakeStack("test-stack") - membershipStack := &generated.Stack{ - ClusterName: "test-stack", - AuthConfig: &generated.AuthConfig{}, - Modules: []*generated.Module{}, // no modules expected → all should be deleted - } - - // Should not panic; errors logged, loop continues - listener.syncModules(ctx, map[string]any{}, stack, membershipStack) -} - -// --- syncStargate error paths --- - -func TestSyncStargate_ApplyFailure(t *testing.T) { - t.Parallel() - ctx := logging.TestingContext() - - k8s := &mockK8SClient{ - applyFn: func(_ context.Context, _ string, _ *unstructured.Unstructured, _ string, _ bool) (*unstructured.Unstructured, error) { - return nil, fmt.Errorf("stargate apply failed") - }, - } - listener := newTestListenerWithMapper(k8s, NewMembershipClientMock(), nil, defaultMapper()) - - stack := newFakeStack("org-stack") - membershipStack := &generated.Stack{ - AuthConfig: &generated.AuthConfig{}, - StargateConfig: &generated.StargateConfig{Enabled: true}, - } - - // Should not panic - listener.syncStargate(ctx, map[string]any{}, stack, membershipStack) -} - -func TestSyncStargate_DeleteFailure(t *testing.T) { - t.Parallel() - ctx := logging.TestingContext() - - k8s := &mockK8SClient{ - ensureNotExistsFn: func(_ context.Context, _, _ string) error { - return fmt.Errorf("delete failed") - }, - } - listener := newTestListenerWithMapper(k8s, NewMembershipClientMock(), nil, defaultMapper()) - - stack := newFakeStack("org-stack") - membershipStack := &generated.Stack{ - StargateConfig: &generated.StargateConfig{Enabled: false}, - } - - // Should not panic - listener.syncStargate(ctx, map[string]any{}, stack, membershipStack) -} - -func TestSyncStargate_NilConfig(t *testing.T) { - t.Parallel() - ctx := logging.TestingContext() - - var deleteCalled bool - k8s := &mockK8SClient{ - ensureNotExistsFn: func(_ context.Context, _, _ string) error { - deleteCalled = true - return nil - }, - } - listener := newTestListenerWithMapper(k8s, NewMembershipClientMock(), nil, defaultMapper()) - - stack := newFakeStack("org-stack") - membershipStack := &generated.Stack{ - StargateConfig: nil, - } - - listener.syncStargate(ctx, map[string]any{}, stack, membershipStack) - assert.True(t, deleteCalled, "should attempt to delete stargate when config is nil") -} - -// --- syncAuthClients error paths --- - -func TestSyncAuthClients_ApplyFailure(t *testing.T) { - t.Parallel() - ctx := logging.TestingContext() - - callCount := 0 - k8s := &mockK8SClient{ - applyFn: func(_ context.Context, _ string, obj *unstructured.Unstructured, _ string, _ bool) (*unstructured.Unstructured, error) { - callCount++ - if callCount == 1 { - return nil, fmt.Errorf("apply failed for first client") - } - result := obj.DeepCopy() - result.SetUID("fake-uid") - return result, nil - }, - } - listener := newTestListenerWithMapper(k8s, NewMembershipClientMock(), nil, defaultMapper()) - - stack := newFakeStack("test-stack") - clients := []*generated.AuthClient{ - {Id: "client1", Public: true}, - {Id: "client2", Public: false}, - } - - // Should not panic; first client fails, second succeeds - listener.syncAuthClients(ctx, map[string]any{}, stack, clients) - assert.Equal(t, 2, callCount, "should attempt apply for both clients") -} - -func TestSyncAuthClients_ListFailure(t *testing.T) { - t.Parallel() - ctx := logging.TestingContext() - - k8s := &mockK8SClient{ - applyFn: func(_ context.Context, _ string, obj *unstructured.Unstructured, _ string, _ bool) (*unstructured.Unstructured, error) { - result := obj.DeepCopy() - result.SetUID("fake-uid") - return result, nil - }, - listFn: func(_ context.Context, _ string, _ labels.Selector) ([]unstructured.Unstructured, error) { - return nil, fmt.Errorf("list forbidden") - }, - } - listener := newTestListenerWithMapper(k8s, NewMembershipClientMock(), nil, defaultMapper()) - - stack := newFakeStack("test-stack") - clients := []*generated.AuthClient{ - {Id: "client1", Public: true}, - } - - // Should not panic; list fails so stale clients are not cleaned up - listener.syncAuthClients(ctx, map[string]any{}, stack, clients) -} - -func TestSyncAuthClients_DeleteStaleFailure(t *testing.T) { - t.Parallel() - ctx := logging.TestingContext() - - var deleteAttempted []string - k8s := &mockK8SClient{ - applyFn: func(_ context.Context, _ string, obj *unstructured.Unstructured, _ string, _ bool) (*unstructured.Unstructured, error) { - result := obj.DeepCopy() - result.SetUID("fake-uid") - return result, nil - }, - listFn: func(_ context.Context, _ string, _ labels.Selector) ([]unstructured.Unstructured, error) { - return []unstructured.Unstructured{ - {Object: map[string]any{"metadata": map[string]any{"name": "test-stack-client1"}}}, - {Object: map[string]any{"metadata": map[string]any{"name": "test-stack-stale"}}}, - }, nil - }, - ensureNotExistsFn: func(_ context.Context, _ string, name string) error { - deleteAttempted = append(deleteAttempted, name) - return fmt.Errorf("delete forbidden") - }, - } - listener := newTestListenerWithMapper(k8s, NewMembershipClientMock(), nil, defaultMapper()) - - stack := newFakeStack("test-stack") - clients := []*generated.AuthClient{ - {Id: "client1", Public: true}, - } - - // Should not panic; stale client deletion fails but is logged - listener.syncAuthClients(ctx, map[string]any{}, stack, clients) - assert.Contains(t, deleteAttempted, "test-stack-stale") -} - -// --- deleteStack error paths --- - -func TestDeleteStack_K8SError(t *testing.T) { - t.Parallel() - ctx := logging.TestingContext() - - k8s := &mockK8SClient{ - deleteFn: func(_ context.Context, _, _ string) error { - return fmt.Errorf("internal server error") - }, - } - mock := NewMembershipClientMock() - listener := newTestListener(k8s, mock, nil) - - // Should not panic; error logged, no message sent to membership - listener.deleteStack(ctx, &generated.DeletedStack{ClusterName: "broken-stack"}) - assert.Empty(t, mock.GetMessages(), "should not send any message on non-NotFound error") -} - -func TestDeleteStack_SendFailure(t *testing.T) { - t.Parallel() - ctx := logging.TestingContext() - - k8s := &mockK8SClient{ - deleteFn: func(_ context.Context, _, _ string) error { - return apierrors.NewNotFound(schema.GroupResource{Group: "formance.com", Resource: "stacks"}, "gone-stack") - }, - } - - // Use a real mock that will track messages, but the test verifies no panic - mock := NewMembershipClientMock() - listener := newTestListener(k8s, mock, nil) - - listener.deleteStack(ctx, &generated.DeletedStack{ClusterName: "gone-stack"}) - require.Len(t, mock.GetMessages(), 1) - assert.NotNil(t, mock.GetMessages()[0].GetStackDeleted()) -} - -// --- disableStack / enableStack error paths --- - -func TestDisableStack_ApplyFailure(t *testing.T) { - t.Parallel() - ctx := logging.TestingContext() - - k8s := &mockK8SClient{ - applyFn: func(_ context.Context, _ string, _ *unstructured.Unstructured, _ string, _ bool) (*unstructured.Unstructured, error) { - return nil, fmt.Errorf("apply failed") - }, - } - listener := newTestListener(k8s, NewMembershipClientMock(), nil) - - // Should not panic - listener.disableStack(ctx, &generated.DisabledStack{ClusterName: "test-stack"}) -} - -func TestEnableStack_ApplyFailure(t *testing.T) { - t.Parallel() - ctx := logging.TestingContext() - - k8s := &mockK8SClient{ - applyFn: func(_ context.Context, _ string, _ *unstructured.Unstructured, _ string, _ bool) (*unstructured.Unstructured, error) { - return nil, fmt.Errorf("apply failed") - }, - } - listener := newTestListener(k8s, NewMembershipClientMock(), nil) - - // Should not panic - listener.enableStack(ctx, &generated.EnabledStack{ClusterName: "test-stack"}) -} - -func TestDisableStack_FieldManager(t *testing.T) { - t.Parallel() - ctx := logging.TestingContext() - - var capturedFieldManager string - var capturedForce bool - k8s := &mockK8SClient{ - applyFn: func(_ context.Context, _ string, _ *unstructured.Unstructured, fm string, force bool) (*unstructured.Unstructured, error) { - capturedFieldManager = fm - capturedForce = force - return &unstructured.Unstructured{}, nil - }, - } - listener := newTestListener(k8s, NewMembershipClientMock(), nil) - - listener.disableStack(ctx, &generated.DisabledStack{ClusterName: "test-stack"}) - assert.Equal(t, fieldManagerDisable, capturedFieldManager) - assert.True(t, capturedForce, "disable should use force=true") -} - -func TestEnableStack_FieldManager(t *testing.T) { - t.Parallel() - ctx := logging.TestingContext() - - var capturedFieldManager string - var capturedForce bool - k8s := &mockK8SClient{ - applyFn: func(_ context.Context, _ string, _ *unstructured.Unstructured, fm string, force bool) (*unstructured.Unstructured, error) { - capturedFieldManager = fm - capturedForce = force - return &unstructured.Unstructured{}, nil - }, - } - listener := newTestListener(k8s, NewMembershipClientMock(), nil) - - listener.enableStack(ctx, &generated.EnabledStack{ClusterName: "test-stack"}) - assert.Equal(t, fieldManagerDisable, capturedFieldManager) - assert.True(t, capturedForce, "enable should use force=true") -} - -// --- createOrUpdate error paths --- - -func TestCreateOrUpdate_ConflictFallsBackToGet(t *testing.T) { - t.Parallel() - ctx := logging.TestingContext() - - existingObj := &unstructured.Unstructured{ - Object: map[string]any{ - "metadata": map[string]any{ - "name": "my-stack", - "annotations": map[string]any{ - "user/custom": "preserved", - }, - }, - }, - } - - k8s := &mockK8SClient{ - applyFn: func(_ context.Context, _ string, _ *unstructured.Unstructured, _ string, _ bool) (*unstructured.Unstructured, error) { - return nil, apierrors.NewConflict(schema.GroupResource{Group: "formance.com", Resource: "stacks"}, "my-stack", fmt.Errorf("field owned by kubectl-edit")) - }, - getFn: func(_ context.Context, _ string, _ string) (*unstructured.Unstructured, error) { - return existingObj, nil - }, - } - - listener := newTestListenerWithMapper(k8s, NewMembershipClientMock(), nil, defaultMapper()) - - result, err := listener.createOrUpdate(ctx, formanceGroupVersion.WithKind("Stack"), "my-stack", "my-stack", nil, map[string]any{ - "spec": map[string]any{"versionsFromFile": "default"}, - }) - - require.NoError(t, err) - assert.Equal(t, "preserved", result.GetAnnotations()["user/custom"]) -} - -func TestCreateOrUpdate_ConflictGetFailure(t *testing.T) { - t.Parallel() - ctx := logging.TestingContext() - - k8s := &mockK8SClient{ - applyFn: func(_ context.Context, _ string, _ *unstructured.Unstructured, _ string, _ bool) (*unstructured.Unstructured, error) { - return nil, apierrors.NewConflict(schema.GroupResource{}, "x", fmt.Errorf("conflict")) - }, - getFn: func(_ context.Context, _ string, _ string) (*unstructured.Unstructured, error) { - return nil, fmt.Errorf("get also failed") - }, - } - - listener := newTestListenerWithMapper(k8s, NewMembershipClientMock(), nil, defaultMapper()) - - _, err := listener.createOrUpdate(ctx, formanceGroupVersion.WithKind("Stack"), "x", "x", nil, map[string]any{ - "spec": map[string]any{}, - }) - - require.Error(t, err) - assert.Contains(t, err.Error(), "getting existing object after conflict") -} - -func TestCreateOrUpdate_NonConflictError(t *testing.T) { - t.Parallel() - ctx := logging.TestingContext() - - k8s := &mockK8SClient{ - applyFn: func(_ context.Context, _ string, _ *unstructured.Unstructured, _ string, _ bool) (*unstructured.Unstructured, error) { - return nil, fmt.Errorf("network error") - }, - } - - listener := newTestListenerWithMapper(k8s, NewMembershipClientMock(), nil, defaultMapper()) - - _, err := listener.createOrUpdate(ctx, formanceGroupVersion.WithKind("Stack"), "x", "x", nil, map[string]any{ - "spec": map[string]any{}, - }) - - require.Error(t, err) - assert.Contains(t, err.Error(), "applying object") -} - -func TestCreateOrUpdate_FieldManagerIsAgent(t *testing.T) { - t.Parallel() - ctx := logging.TestingContext() - - var capturedFM string - var capturedForce bool - k8s := &mockK8SClient{ - applyFn: func(_ context.Context, _ string, obj *unstructured.Unstructured, fm string, force bool) (*unstructured.Unstructured, error) { - capturedFM = fm - capturedForce = force - result := obj.DeepCopy() - result.SetUID("fake-uid") - return result, nil - }, - } - - listener := newTestListenerWithMapper(k8s, NewMembershipClientMock(), nil, defaultMapper()) - - _, err := listener.createOrUpdate(ctx, formanceGroupVersion.WithKind("Stack"), "x", "x", nil, map[string]any{ - "spec": map[string]any{}, - }) - - require.NoError(t, err) - assert.Equal(t, fieldManagerAgent, capturedFM) - assert.False(t, capturedForce, "createOrUpdate should use force=false") -} - -// --- informer_versions (entirely untested until now) --- - -func TestVersionsEventHandler_Add(t *testing.T) { - t.Parallel() - - mock := NewMembershipClientMock() - handler := VersionsEventHandler(logging.Testing(), mock) - - version := &unstructured.Unstructured{ - Object: map[string]any{ - "metadata": map[string]any{ - "name": "v2.0", - "annotations": map[string]any{ - "formance.com/deprecated": "true", - }, - }, - "spec": map[string]any{ - "ledger": "v2.0.0", - "payments": "v1.5.0", - }, - }, - } - - handler.OnAdd(version, false) - - require.Len(t, mock.GetMessages(), 1) - msg := mock.GetMessages()[0].GetAddedVersion() - require.NotNil(t, msg) - assert.Equal(t, "v2.0", msg.Name) - assert.Equal(t, "v2.0.0", msg.Versions["ledger"]) - assert.Equal(t, "v1.5.0", msg.Versions["payments"]) - assert.True(t, msg.Deprecated) -} - -func TestVersionsEventHandler_Add_NoSpec(t *testing.T) { - t.Parallel() - - mock := NewMembershipClientMock() - handler := VersionsEventHandler(logging.Testing(), mock) - - version := &unstructured.Unstructured{ - Object: map[string]any{ - "metadata": map[string]any{ - "name": "empty", - }, - }, - } - - handler.OnAdd(version, false) - - require.Len(t, mock.GetMessages(), 1) - msg := mock.GetMessages()[0].GetAddedVersion() - require.NotNil(t, msg) - assert.Nil(t, msg.Versions) - assert.False(t, msg.Deprecated) -} - -func TestVersionsEventHandler_Update_Changed(t *testing.T) { - t.Parallel() - - mock := NewMembershipClientMock() - handler := VersionsEventHandler(logging.Testing(), mock) - - oldVersion := &unstructured.Unstructured{ - Object: map[string]any{ - "metadata": map[string]any{"name": "v1"}, - "spec": map[string]any{"ledger": "v1.0.0"}, - }, - } - newVersion := &unstructured.Unstructured{ - Object: map[string]any{ - "metadata": map[string]any{"name": "v1"}, - "spec": map[string]any{"ledger": "v1.1.0"}, - }, - } - - handler.OnUpdate(oldVersion, newVersion) - - require.Len(t, mock.GetMessages(), 1) - msg := mock.GetMessages()[0].GetUpdatedVersion() - require.NotNil(t, msg) - assert.Equal(t, "v1.1.0", msg.Versions["ledger"]) -} - -func TestVersionsEventHandler_Update_Unchanged(t *testing.T) { - t.Parallel() - - mock := NewMembershipClientMock() - handler := VersionsEventHandler(logging.Testing(), mock) - - obj := &unstructured.Unstructured{ - Object: map[string]any{ - "metadata": map[string]any{"name": "v1"}, - "spec": map[string]any{"ledger": "v1.0.0"}, - }, - } - - handler.OnUpdate(obj, obj) - - assert.Empty(t, mock.GetMessages(), "should not send message when spec unchanged") -} - -func TestVersionsEventHandler_Delete(t *testing.T) { - t.Parallel() - - mock := NewMembershipClientMock() - handler := VersionsEventHandler(logging.Testing(), mock) - - version := &unstructured.Unstructured{ - Object: map[string]any{ - "metadata": map[string]any{"name": "v1"}, - }, - } - - handler.OnDelete(version) - - require.Len(t, mock.GetMessages(), 1) - msg := mock.GetMessages()[0].GetDeletedVersion() - require.NotNil(t, msg) - assert.Equal(t, "v1", msg.Name) -} - -// --- extractVersionsSpec --- - -func TestExtractVersionsSpec(t *testing.T) { - t.Parallel() - - tests := []struct { - name string - obj *unstructured.Unstructured - expect map[string]string - }{ - { - name: "normal spec", - obj: &unstructured.Unstructured{ - Object: map[string]any{ - "spec": map[string]any{ - "ledger": "v1.0", "payments": "v2.0", - }, - }, - }, - expect: map[string]string{"ledger": "v1.0", "payments": "v2.0"}, - }, - { - name: "nil spec", - obj: &unstructured.Unstructured{ - Object: map[string]any{}, - }, - expect: nil, - }, - { - name: "non-string values ignored", - obj: &unstructured.Unstructured{ - Object: map[string]any{ - "spec": map[string]any{ - "ledger": "v1.0", - "count": int64(42), - }, - }, - }, - expect: map[string]string{"ledger": "v1.0"}, - }, - } - - for _, tc := range tests { - t.Run(tc.name, func(t *testing.T) { - t.Parallel() - result := extractVersionsSpec(tc.obj) - assert.Equal(t, tc.expect, result) - }) - } -} - -// --- generateMetadata --- - -func TestGenerateMetadata(t *testing.T) { - t.Parallel() - - listener := newTestListener(&mockK8SClient{}, NewMembershipClientMock(), nil) - - t.Run("with labels and annotations", func(t *testing.T) { - stack := &generated.Stack{ - AdditionalLabels: map[string]string{"env": "prod", "team": "infra"}, - AdditionalAnnotations: map[string]string{"note": "important"}, - } - md := listener.generateMetadata(stack) - - labels := md["labels"].(map[string]any) - assert.Equal(t, "prod", labels["formance.com/env"]) - assert.Equal(t, "infra", labels["formance.com/team"]) - - annotations := md["annotations"].(map[string]any) - assert.Equal(t, "important", annotations["formance.com/note"]) - }) - - t.Run("nil labels and annotations", func(t *testing.T) { - stack := &generated.Stack{} - md := listener.generateMetadata(stack) - - labels := md["labels"].(map[string]any) - assert.Empty(t, labels) - - annotations := md["annotations"].(map[string]any) - assert.Empty(t, annotations) - }) -} diff --git a/internal/k8s_client_mock_test.go b/internal/k8s_client_mock_test.go deleted file mode 100644 index 5d61216..0000000 --- a/internal/k8s_client_mock_test.go +++ /dev/null @@ -1,63 +0,0 @@ -package internal - -import ( - "context" - - "k8s.io/apimachinery/pkg/apis/meta/v1/unstructured" - "k8s.io/apimachinery/pkg/labels" -) - -// mockK8SClient is a configurable mock for unit-testing error paths. -// Each field is a function that, when set, overrides the default (no-op) behaviour. -type mockK8SClient struct { - getFn func(ctx context.Context, resource, name string) (*unstructured.Unstructured, error) - applyFn func(ctx context.Context, resource string, obj *unstructured.Unstructured, fieldManager string, force bool) (*unstructured.Unstructured, error) - deleteFn func(ctx context.Context, resource, name string) error - ensureNotExistsFn func(ctx context.Context, resource, name string) error - ensureNotExistsBySelectorFn func(ctx context.Context, resource string, selector labels.Selector) error - listFn func(ctx context.Context, resource string, selector labels.Selector) ([]unstructured.Unstructured, error) -} - -func (m *mockK8SClient) Get(ctx context.Context, resource, name string) (*unstructured.Unstructured, error) { - if m.getFn != nil { - return m.getFn(ctx, resource, name) - } - return &unstructured.Unstructured{}, nil -} - -func (m *mockK8SClient) Apply(ctx context.Context, resource string, obj *unstructured.Unstructured, fieldManager string, force bool) (*unstructured.Unstructured, error) { - if m.applyFn != nil { - return m.applyFn(ctx, resource, obj, fieldManager, force) - } - result := obj.DeepCopy() - result.SetUID("fake-uid") - return result, nil -} - -func (m *mockK8SClient) Delete(ctx context.Context, resource, name string) error { - if m.deleteFn != nil { - return m.deleteFn(ctx, resource, name) - } - return nil -} - -func (m *mockK8SClient) EnsureNotExists(ctx context.Context, resource, name string) error { - if m.ensureNotExistsFn != nil { - return m.ensureNotExistsFn(ctx, resource, name) - } - return nil -} - -func (m *mockK8SClient) EnsureNotExistsBySelector(ctx context.Context, resource string, selector labels.Selector) error { - if m.ensureNotExistsBySelectorFn != nil { - return m.ensureNotExistsBySelectorFn(ctx, resource, selector) - } - return nil -} - -func (m *mockK8SClient) List(ctx context.Context, resource string, selector labels.Selector) ([]unstructured.Unstructured, error) { - if m.listFn != nil { - return m.listFn(ctx, resource, selector) - } - return nil, nil -} diff --git a/tests/membership_listener_test.go b/tests/membership_listener_test.go index caacffa..fb3339b 100644 --- a/tests/membership_listener_test.go +++ b/tests/membership_listener_test.go @@ -227,6 +227,62 @@ var _ = Describe("Membership listener", func() { Expect(u).To(TargetStack(stack)) } }) + It("Should default versionsFromFile to 'default' when empty", func() { + s := &unstructured.Unstructured{} + Expect(LoadResource("Stacks", membershipStack.ClusterName, s)).To(Succeed()) + versionsFromFile, _, _ := unstructured.NestedString(s.Object, "spec", "versionsFromFile") + Expect(versionsFromFile).To(Equal("default")) + }) + It("Should use the SSA field manager 'formance-agent'", func() { + s := &unstructured.Unstructured{} + Expect(LoadResource("Stacks", membershipStack.ClusterName, s)).To(Succeed()) + + managedFields := s.GetManagedFields() + found := false + for _, mf := range managedFields { + if mf.Manager == "formance-agent" { + found = true + break + } + } + Expect(found).To(BeTrue(), "expected managedFields to contain 'formance-agent'") + }) + It("Should be idempotent on resync", func() { + s1 := &unstructured.Unstructured{} + Expect(LoadResource("Stacks", membershipStack.ClusterName, s1)).To(Succeed()) + rv1 := s1.GetResourceVersion() + + membershipClient.Orders() <- &generated.Order{ + Message: &generated.Order_ExistingStack{ + ExistingStack: membershipStack, + }, + } + + // Give the agent time to process the order + Consistently(func(g Gomega) { + s2 := &unstructured.Unstructured{} + g.Expect(LoadResource("Stacks", membershipStack.ClusterName, s2)).To(Succeed()) + g.Expect(s2.GetResourceVersion()).To(Equal(rv1)) + }, "2s", "200ms").Should(Succeed()) + }) + When("versions is set explicitly", func() { + BeforeEach(func() { + membershipStack.Versions = "v1.2.3" + membershipClient.Orders() <- &generated.Order{ + Message: &generated.Order_ExistingStack{ + ExistingStack: membershipStack, + }, + } + }) + It("Should use the explicit version", func() { + Eventually(func(g Gomega) { + s := &unstructured.Unstructured{} + g.Expect(LoadResource("Stacks", membershipStack.ClusterName, s)).To(Succeed()) + versionsFromFile, _, _ := unstructured.NestedString(s.Object, "spec", "versionsFromFile") + g.Expect(versionsFromFile).To(Equal("v1.2.3")) + }).Should(Succeed()) + }) + }) When("a user manually edits a resource managed by the agent", func() { It("Should preserve user-added annotations after resync", func() { userPatch, err := json.Marshal(map[string]any{ diff --git a/tests/versions_informer_test.go b/tests/versions_informer_test.go new file mode 100644 index 0000000..0677d0e --- /dev/null +++ b/tests/versions_informer_test.go @@ -0,0 +1,192 @@ +package tests + +import ( + "context" + "encoding/json" + "time" + + "github.com/formancehq/go-libs/v2/logging" + "github.com/formancehq/stack/components/agent/internal" + "github.com/google/uuid" + . "github.com/onsi/ginkgo/v2" + . "github.com/onsi/gomega" + "k8s.io/apimachinery/pkg/apis/meta/v1/unstructured" + "k8s.io/apimachinery/pkg/types" + "k8s.io/client-go/dynamic" +) + +var _ = Describe("Versions informer", func() { + var ( + membershipClientMock *internal.MembershipClientMock + startListener func() + ) + BeforeEach(func() { + membershipClientMock = internal.NewMembershipClientMock() + dynamicClient, err := dynamic.NewForConfig(restConfig) + Expect(err).To(Succeed()) + + factory := internal.NewDynamicSharedInformerFactory(dynamicClient, 5*time.Minute) + Expect(internal.CreateVersionsInformer(factory, logging.Testing(), membershipClientMock)).To(Succeed()) + startListener = func() { + stopCh := make(chan struct{}) + factory.Start(stopCh) + DeferCleanup(func() { + close(stopCh) + }) + } + }) + When("a Versions resource is created", func() { + var version *unstructured.Unstructured + BeforeEach(func() { + version = &unstructured.Unstructured{ + Object: map[string]interface{}{ + "apiVersion": formanceGroupVersion.String(), + "kind": "Versions", + "metadata": map[string]interface{}{ + "name": uuid.NewString(), + "annotations": map[string]interface{}{ + "formance.com/deprecated": "true", + }, + }, + "spec": map[string]interface{}{ + "ledger": "v2.0.0", + "payments": "v1.5.0", + }, + }, + } + Expect(k8sClient.Post(). + Resource("Versions"). + Body(version). + Do(context.Background()). + Into(version)).To(Succeed()) + + startListener() + + DeferCleanup(func() { + k8sClient.Delete(). + Resource("Versions"). + Name(version.GetName()). + Do(context.Background()) + }) + }) + It("Should send AddedVersion with spec and deprecated flag", func() { + Eventually(func(g Gomega) { + for _, message := range membershipClientMock.GetMessages() { + if msg := message.GetAddedVersion(); msg != nil && msg.Name == version.GetName() { + g.Expect(msg.Versions["ledger"]).To(Equal("v2.0.0")) + g.Expect(msg.Versions["payments"]).To(Equal("v1.5.0")) + g.Expect(msg.Deprecated).To(BeTrue()) + return + } + } + g.Expect(false).To(BeTrue(), "AddedVersion message not found") + }).Should(Succeed()) + }) + When("the spec is updated", func() { + BeforeEach(func() { + // Wait for AddedVersion first + Eventually(func() bool { + for _, message := range membershipClientMock.GetMessages() { + if msg := message.GetAddedVersion(); msg != nil && msg.Name == version.GetName() { + return true + } + } + return false + }).Should(BeTrue()) + + patch, err := json.Marshal(map[string]any{ + "spec": map[string]any{ + "ledger": "v2.1.0", + "payments": "v1.5.0", + }, + }) + Expect(err).To(BeNil()) + Expect(k8sClient.Patch(types.MergePatchType). + Resource("Versions"). + Name(version.GetName()). + Body(patch). + Do(context.Background()). + Error()).To(Succeed()) + }) + It("Should send UpdatedVersion", func() { + Eventually(func(g Gomega) { + for _, message := range membershipClientMock.GetMessages() { + if msg := message.GetUpdatedVersion(); msg != nil && msg.Name == version.GetName() { + g.Expect(msg.Versions["ledger"]).To(Equal("v2.1.0")) + return + } + } + g.Expect(false).To(BeTrue(), "UpdatedVersion message not found") + }).Should(Succeed()) + }) + }) + When("the resource is deleted", func() { + BeforeEach(func() { + // Wait for AddedVersion first + Eventually(func() bool { + for _, message := range membershipClientMock.GetMessages() { + if msg := message.GetAddedVersion(); msg != nil && msg.Name == version.GetName() { + return true + } + } + return false + }).Should(BeTrue()) + + Expect(k8sClient.Delete(). + Resource("Versions"). + Name(version.GetName()). + Do(context.Background()).Error()).To(Succeed()) + }) + It("Should send DeletedVersion", func() { + Eventually(func(g Gomega) { + for _, message := range membershipClientMock.GetMessages() { + if msg := message.GetDeletedVersion(); msg != nil && msg.Name == version.GetName() { + return + } + } + g.Expect(false).To(BeTrue(), "DeletedVersion message not found") + }).Should(Succeed()) + }) + }) + }) + When("a Versions resource has no spec", func() { + var version *unstructured.Unstructured + BeforeEach(func() { + version = &unstructured.Unstructured{ + Object: map[string]interface{}{ + "apiVersion": formanceGroupVersion.String(), + "kind": "Versions", + "metadata": map[string]interface{}{ + "name": uuid.NewString(), + }, + }, + } + Expect(k8sClient.Post(). + Resource("Versions"). + Body(version). + Do(context.Background()). + Into(version)).To(Succeed()) + + startListener() + + DeferCleanup(func() { + k8sClient.Delete(). + Resource("Versions"). + Name(version.GetName()). + Do(context.Background()) + }) + }) + It("Should send AddedVersion with nil versions", func() { + Eventually(func(g Gomega) { + for _, message := range membershipClientMock.GetMessages() { + if msg := message.GetAddedVersion(); msg != nil && msg.Name == version.GetName() { + g.Expect(msg.Versions).To(BeEmpty()) + g.Expect(msg.Deprecated).To(BeFalse()) + return + } + } + g.Expect(false).To(BeTrue(), "AddedVersion message not found") + }).Should(Succeed()) + }) + }) +})