diff --git a/go.mod b/go.mod index b4bbc850..d848ac44 100644 --- a/go.mod +++ b/go.mod @@ -9,6 +9,7 @@ require ( github.com/go-playground/validator/v10 v10.30.3 github.com/go-viper/mapstructure/v2 v2.5.0 github.com/mitchellh/copystructure v1.2.0 + github.com/openshift-hyperfleet/hyperfleet-applier v0.0.0-20260827133207-94b7a4d56697 github.com/openshift-hyperfleet/hyperfleet-broker v1.1.1 github.com/openshift-online/maestro v0.0.0-20260202062555-48b47506a254 github.com/openshift-online/ocm-sdk-go v0.1.510 @@ -128,6 +129,7 @@ require ( github.com/oklog/ulid v1.3.1 // indirect github.com/opencontainers/go-digest v1.0.0 // indirect github.com/opencontainers/image-spec v1.1.1 // indirect + github.com/openshift-hyperfleet/hyperfleet-logger v0.0.0-20260811173525-c9f9e282d029 // indirect github.com/pelletier/go-toml/v2 v2.4.2 // indirect github.com/pkg/errors v0.9.1 // indirect github.com/pmezard/go-difflib v1.0.1-0.20181226105442-5d4384ee4fb2 // indirect @@ -180,7 +182,7 @@ require ( gopkg.in/yaml.v2 v2.4.0 // indirect k8s.io/api v0.37.0 // indirect k8s.io/klog/v2 v2.140.0 // indirect - k8s.io/kube-openapi v0.0.0-20260721132016-d427ff9ee9ad // indirect + k8s.io/kube-openapi v0.0.0-20260821135717-be32def86098 // indirect k8s.io/utils v0.0.0-20260707023825-cf1189d6abe3 // indirect sigs.k8s.io/json v0.0.0-20250730193827-2d320260d730 // indirect sigs.k8s.io/randfill v1.0.0 // indirect diff --git a/go.sum b/go.sum index 0df42505..e1b729ef 100644 --- a/go.sum +++ b/go.sum @@ -273,8 +273,12 @@ github.com/opencontainers/go-digest v1.0.0 h1:apOUWs51W5PlhuyGyz9FCeeBIOUDA/6nW8 github.com/opencontainers/go-digest v1.0.0/go.mod h1:0JzlMkj0TRzQZfJkVvzbP0HBR3IKzErnv2BNG4W4MAM= github.com/opencontainers/image-spec v1.1.1 h1:y0fUlFfIZhPF1W537XOLg0/fcx6zcHCJwooC2xJA040= github.com/opencontainers/image-spec v1.1.1/go.mod h1:qpqAh3Dmcf36wStyyWU+kCeDgrGnAve2nCC8+7h8Q0M= +github.com/openshift-hyperfleet/hyperfleet-applier v0.0.0-20260827133207-94b7a4d56697 h1:9DFSlPuXlMY4zshztT4/MEnY6QAgu0WoVY1KFtVE9fg= +github.com/openshift-hyperfleet/hyperfleet-applier v0.0.0-20260827133207-94b7a4d56697/go.mod h1:tDZnkLwpVO5ksfPm3NH9dy5Yl0ZBSfbtOA3hj2wrmHc= github.com/openshift-hyperfleet/hyperfleet-broker v1.1.1 h1:3zbpNuFL+OEvKl6a/KJAlHFcYR4QqQBoGkf35higypU= github.com/openshift-hyperfleet/hyperfleet-broker v1.1.1/go.mod h1:E7Br4NnsaTTfWR2fEqHAtvFXUAgzFpksF+G5qTBMmy0= +github.com/openshift-hyperfleet/hyperfleet-logger v0.0.0-20260811173525-c9f9e282d029 h1:c3GdD3EUdR9lRjNot+aBIxVbAXbOSTtiQpWVTHnFHaU= +github.com/openshift-hyperfleet/hyperfleet-logger v0.0.0-20260811173525-c9f9e282d029/go.mod h1:5Nh2IMS2MouehZ4UiRXFMjhuwdys2DX26J4vArmXe6Y= github.com/openshift-online/maestro v0.0.0-20260202062555-48b47506a254 h1:v/jYqdzZpzB/bscVpajlbcKgCNeV4tx4fkm5m2JR8Ug= github.com/openshift-online/maestro v0.0.0-20260202062555-48b47506a254/go.mod h1:cyeif610uObNrbcyn5s1fZg7OWseVjaMAqgrEDA2Aec= github.com/openshift-online/ocm-sdk-go v0.1.510 h1:oPwgGHPi6LeY+RANXUhJwVvkY4yvddyRqHHThHZZRyM= @@ -508,16 +512,16 @@ honnef.co/go/tools v0.0.0-20190102054323-c2f93a96b099/go.mod h1:rf3lG4BRIbNafJWh honnef.co/go/tools v0.0.0-20190523083050-ea95bdfd59fc/go.mod h1:rf3lG4BRIbNafJWhAfAdb/ePZxsR/4RtNHQocxwk9r4= k8s.io/api v0.37.0 h1:Z//Vj9N7RA/yS2sDmxyeo7h+RR4zbUrd2vrd3Z0TbB4= k8s.io/api v0.37.0/go.mod h1:LKXgcJWMc+f4OLbP5SFR8rulEg07zZhpi/zMULiBImk= -k8s.io/apiextensions-apiserver v0.36.0 h1:Wt7E8J+VBCbj4FjiBfDTK/neXDDjyJVJc7xfuOHImZ0= -k8s.io/apiextensions-apiserver v0.36.0/go.mod h1:kGDjH0msuiIB3tgsYRV0kS9GqpMYMUsQ3GHv7TApyug= +k8s.io/apiextensions-apiserver v0.36.4 h1:SfvCVt+4CqKWvzuVytYDT5g9hyb9MztoiYELIkPVrFc= +k8s.io/apiextensions-apiserver v0.36.4/go.mod h1:JT9V2Ju7ys1FY4zbSpmX9XOvKB3/BwsODc4hFQEa+Xo= k8s.io/apimachinery v0.37.0 h1:Np2AbDtf8x6RDHiD8T9LbKJ9gaegeVNa8yNm5FuGKm0= k8s.io/apimachinery v0.37.0/go.mod h1:RN3nhprFSCxOi5Selxd7oMTXOe/c+ZbcE7Im+TS2zkE= k8s.io/client-go v0.37.0 h1:nsN31fy8wBySuZ+QRnKmrjRSQLOG2rvoGN0tKd12zhQ= k8s.io/client-go v0.37.0/go.mod h1:FcGqw+Ll/gNQiq+nPGY1Oyt9y7SgDh1d3MW3RFDEbn0= k8s.io/klog/v2 v2.140.0 h1:Tf+J3AH7xnUzZyVVXhTgGhEKnFqye14aadWv7bzXdzc= k8s.io/klog/v2 v2.140.0/go.mod h1:o+/RWfJ6PwpnFn7OyAG3QnO47BFsymfEfrz6XyYSSp0= -k8s.io/kube-openapi v0.0.0-20260721132016-d427ff9ee9ad h1:oXImqH8mQNk7PmvzKhmN3ddJoY6OnyM225MXwGHPm0A= -k8s.io/kube-openapi v0.0.0-20260721132016-d427ff9ee9ad/go.mod h1:0/mqHCVhlumdJ3BhCfnjSZQE037nAhNodh1/hK0T8/I= +k8s.io/kube-openapi v0.0.0-20260821135717-be32def86098 h1:z5+pcu1jTyKK5mNTe2/+x+U6Uuv9jRVOJQLaBJJMpeI= +k8s.io/kube-openapi v0.0.0-20260821135717-be32def86098/go.mod h1:0/mqHCVhlumdJ3BhCfnjSZQE037nAhNodh1/hK0T8/I= k8s.io/utils v0.0.0-20260707023825-cf1189d6abe3 h1:jVkFFVfXdXP74B/zbO3hM3hpSFD0xvhQ5U686DPurkE= k8s.io/utils v0.0.0-20260707023825-cf1189d6abe3/go.mod h1:M2s5JB1lIYP3jzZdorPLHXIPJzt9vv2muW5a6L9DtNM= open-cluster-management.io/api v1.3.0 h1:Q3miH38BE3N5+PesHQ0kcFi5nhX5350m7OJWapZcVqY= diff --git a/internal/desireclient/apply.go b/internal/desireclient/apply.go new file mode 100644 index 00000000..11490284 --- /dev/null +++ b/internal/desireclient/apply.go @@ -0,0 +1,173 @@ +package desireclient + +import ( + "context" + "encoding/json" + "errors" + "fmt" + + "github.com/openshift-hyperfleet/hyperfleet-adapter/internal/manifest" + "github.com/openshift-hyperfleet/hyperfleet-adapter/internal/transportclient" + "github.com/openshift-hyperfleet/hyperfleet-applier/pkg/desire" +) + +// ApplyResource implements transportclient.TransportClient. It upserts an +// apply desire for the rendered manifest and auto-creates its paired read +// desire so the applied resource becomes visible to discovery. A read-desire +// pairing failure is returned as an error: an apply desire without its read +// desire is permanently invisible to discovery. +func (c *Client) ApplyResource( + ctx context.Context, + manifestBytes []byte, + opts *transportclient.ApplyOptions, + target transportclient.TransportContext, +) (*transportclient.ApplyResult, error) { + if len(manifestBytes) == 0 { + return nil, fmt.Errorf("desireclient: manifest bytes cannot be empty") + } + + tc, err := resolveTransportContext(target) + if err != nil { + return nil, err + } + + obj, err := parseToUnstructured(manifestBytes) + if err != nil { + return nil, fmt.Errorf("desireclient: failed to parse manifest: %w", err) + } + + // The store's ApplySpec.KubeContent must be valid JSON, but manifestBytes + // may have been YAML (parseToUnstructured accepts both). Re-marshal the + // parsed object rather than storing the original bytes verbatim. + kubeContent, err := json.Marshal(obj.Object) + if err != nil { + return nil, fmt.Errorf("desireclient: failed to marshal manifest to JSON: %w", err) + } + + gvk := obj.GroupVersionKind() + namespace, name := obj.GetNamespace(), obj.GetName() + + readID, err := buildIdentity(tc, desire.TypeRead, gvk, namespace, name) + if err != nil { + return nil, err + } + + applyID, err := buildIdentity(tc, desire.TypeApply, gvk, namespace, name) + if err != nil { + return nil, err + } + + existing, err := c.store.GetApplyDesire(ctx, applyID) + if err != nil && !errors.Is(err, desire.ErrNotFound) { + return nil, fmt.Errorf("desireclient: failed to get apply desire for %s/%s: %w", applyID.Namespace, applyID.Name, err) + } + exists := err == nil + + newGen := manifest.GetGenerationFromUnstructured(obj) + var existingGen int64 + if exists { + existingGen = generationFromKubeContent(existing.Spec.KubeContent) + } + + decision := manifest.CompareGenerations(newGen, existingGen, exists) + result := &transportclient.ApplyResult{Operation: decision.Operation, Reason: decision.Reason} + + // 2. Pairing self-heals on every path, including skip: an externally + // deleted read desire is re-created on the next event. But the mirror + // only recreate to a new TargetVersion when this call is actually + // writing new apply content — on skip, ApplyDesire.Spec.KubeContent + // stays at the old API version, so the mirror must too. They move + // together or not at all. + recreate := decision.Operation != manifest.OperationSkip + if err = c.ensureReadDesire(ctx, readID, gvk.Version, recreate); err != nil { + return nil, fmt.Errorf( + "desireclient: failed to create paired read desire for %s/%s: %w", gvk.Kind, name, err) + } + + switch decision.Operation { + case manifest.OperationCreate: + if _, err = c.store.CreateApplyDesire(ctx, desire.ApplyDesire{ + Identity: applyID, + Owner: c.owner, + Spec: desire.ApplySpec{KubeContent: kubeContent}, + }); err != nil { + c.logApplyError(ctx, applyID, err) + return nil, fmt.Errorf("desireclient: failed to create apply desire for %s/%s: %w", namespace, name, err) + } + case manifest.OperationUpdate: + if _, err = c.store.UpdateApplyDesireSpec( + ctx, applyID, desire.ApplySpec{KubeContent: kubeContent}, c.owner, existing.Version, + ); err != nil { + c.logApplyError(ctx, applyID, err) + return nil, fmt.Errorf( + "desireclient: failed to update apply desire for %s/%s: %w", applyID.Namespace, applyID.Name, err) + } + case manifest.OperationSkip: + // Nothing to do. + default: + return nil, fmt.Errorf("desireclient: unexpected apply decision operation %q", decision.Operation) + } + + c.log.Debugf(ctx, "ApplyResource %s/%s: operation=%s reason=%s", + applyID.Namespace, applyID.Name, result.Operation, result.Reason) + return result, nil +} + +func (c *Client) logApplyError(ctx context.Context, id desire.Identity, err error) { + switch { + case errors.Is(err, desire.ErrOwnerConflict): + c.log.Errorf(ctx, "ApplyResource %s/%s: operation failed with ErrOwnerConflict", id.Namespace, id.Name) + case errors.Is(err, desire.ErrDeletePending): + c.log.Warnf(ctx, "ApplyResource %s/%s: operation failed with ErrDeletePending", id.Namespace, id.Name) + case errors.Is(err, desire.ErrVersionConflict): + c.log.Warnf(ctx, "ApplyResource %s/%s: operation failed with ErrVersionConflict", id.Namespace, id.Name) + default: + c.log.Errorf(ctx, "ApplyResource %s/%s: operation failed with unexpected error: %v", id.Namespace, id.Name, err) + } +} + +// ensureReadDesire creates the paired read desire if it doesn't already +// exist. ErrAlreadyExists is a no-op success. +func (c *Client) ensureReadDesire(ctx context.Context, id desire.Identity, targetVersion string, recreate bool) error { + _, err := c.store.CreateReadDesire(ctx, desire.ReadDesire{ + Identity: id, + Owner: c.owner, + TargetVersion: targetVersion, + }) + if err == nil { + return nil + } + + if !errors.Is(err, desire.ErrAlreadyExists) { + return fmt.Errorf("desireclient: failed to create read desire for %s/%s: %w", id.Namespace, id.Name, err) + } + + if !recreate { + return nil + } + + des, err := c.store.GetReadDesire(ctx, id) + if err != nil { + return fmt.Errorf("desireclient: failed to get existing read desire for %s/%s: %w", id.Namespace, id.Name, err) + } + + if des.TargetVersion == targetVersion { + return nil + } + + err = c.store.DeleteReadDesire(ctx, id, c.owner, des.Version) + if err != nil { + return fmt.Errorf("desireclient: failed to delete existing read desire for %s/%s: %w", id.Namespace, id.Name, err) + } + + if _, err := c.store.CreateReadDesire(ctx, desire.ReadDesire{ + Identity: id, + Owner: c.owner, + TargetVersion: targetVersion, + }); err != nil { + return fmt.Errorf("desireclient: failed to create read desire for %s/%s: %w", id.Namespace, id.Name, err) + } + + return nil + +} diff --git a/internal/desireclient/apply_test.go b/internal/desireclient/apply_test.go new file mode 100644 index 00000000..c9c5411e --- /dev/null +++ b/internal/desireclient/apply_test.go @@ -0,0 +1,238 @@ +package desireclient + +import ( + "context" + "encoding/json" + "errors" + "testing" + + "github.com/openshift-hyperfleet/hyperfleet-adapter/internal/manifest" + "github.com/openshift-hyperfleet/hyperfleet-applier/pkg/desire" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +func TestApplyResource_CreatesApplyAndReadDesire(t *testing.T) { + ctx := context.Background() + store := newMemoryStore() + c := newTestClient(store) + + result, err := c.ApplyResource(ctx, configMapManifest(1), nil, testTransportContext()) + require.NoError(t, err) + assert.Equal(t, manifest.OperationCreate, result.Operation) + + applyID := desire.Identity{ + ManagementCluster: testManagementCluster, Type: desire.TypeApply, + Resource: testResource, Namespace: testNamespace, Name: testName, + } + applied, err := store.GetApplyDesire(ctx, applyID) + require.NoError(t, err) + assert.Equal(t, testOwner, applied.Owner) + + readID := applyID + readID.Type = desire.TypeRead + readDesire, err := store.GetReadDesire(ctx, readID) + require.NoError(t, err, "paired read desire must be auto-created") + assert.Equal(t, "v1", readDesire.TargetVersion) + assert.Equal(t, testOwner, readDesire.Owner) +} + +func TestApplyResource_UpdatesWithCASVersion(t *testing.T) { + ctx := context.Background() + store := newMemoryStore() + c := newTestClient(store) + + _, err := c.ApplyResource(ctx, configMapManifest(1), nil, testTransportContext()) + require.NoError(t, err) + + applyID := desire.Identity{ + ManagementCluster: testManagementCluster, Type: desire.TypeApply, + Resource: testResource, Namespace: testNamespace, Name: testName, + } + before, err := store.GetApplyDesire(ctx, applyID) + require.NoError(t, err) + require.Equal(t, int64(1), before.Version) + + result, err := c.ApplyResource(ctx, configMapManifest(2), nil, testTransportContext()) + require.NoError(t, err) + assert.Equal(t, manifest.OperationUpdate, result.Operation) + + after, err := store.GetApplyDesire(ctx, applyID) + require.NoError(t, err) + assert.Equal(t, int64(2), after.Version, "UpdateApplyDesireSpec must have used the fetched CAS version") + assert.Contains(t, string(after.Spec.KubeContent), `"key":"value"`) +} + +func TestApplyResource_AcceptsYAMLManifest(t *testing.T) { + ctx := context.Background() + store := newMemoryStore() + c := newTestClient(store) + + yamlManifest := []byte(` +apiVersion: v1 +kind: ConfigMap +metadata: + name: my-config + namespace: default + annotations: + hyperfleet.io/generation: "1" +data: + key: value +`) + + result, err := c.ApplyResource(ctx, yamlManifest, nil, testTransportContext()) + require.NoError(t, err) + assert.Equal(t, manifest.OperationCreate, result.Operation) + + applyID := desire.Identity{ + ManagementCluster: testManagementCluster, Type: desire.TypeApply, + Resource: testResource, Namespace: testNamespace, Name: testName, + } + applied, err := store.GetApplyDesire(ctx, applyID) + require.NoError(t, err) + // The store requires KubeContent to be valid JSON; a YAML input must be + // normalized before being persisted, not stored verbatim. + assert.True(t, json.Valid(applied.Spec.KubeContent), + "KubeContent must be valid JSON even when the input manifest was YAML") + assert.Contains(t, string(applied.Spec.KubeContent), `"key":"value"`) +} + +func TestApplyResource_SkipsWhenGenerationUnchanged(t *testing.T) { + ctx := context.Background() + store := newMemoryStore() + c := newTestClient(store) + + _, err := c.ApplyResource(ctx, configMapManifest(1), nil, testTransportContext()) + require.NoError(t, err) + + result, err := c.ApplyResource(ctx, configMapManifest(1), nil, testTransportContext()) + require.NoError(t, err) + assert.Equal(t, manifest.OperationSkip, result.Operation) +} + +func TestApplyResource_ReadDesirePairingFailureIsError(t *testing.T) { + ctx := context.Background() + store := &failingReadDesireStore{SpecStore: newMemoryStore()} + c := newTestClient(store) + + _, err := c.ApplyResource(ctx, configMapManifest(1), nil, testTransportContext()) + require.Error(t, err, "a read-desire pairing failure must surface as an error, not a warning") + assert.Contains(t, err.Error(), "paired read desire") +} + +func TestApplyResource_RequiresTransportContext(t *testing.T) { + ctx := context.Background() + c := newTestClient(newMemoryStore()) + + _, err := c.ApplyResource(ctx, configMapManifest(1), nil, nil) + require.Error(t, err) +} + +func TestApplyResource_EmptyManifestIsError(t *testing.T) { + ctx := context.Background() + c := newTestClient(newMemoryStore()) + + _, err := c.ApplyResource(ctx, nil, nil, testTransportContext()) + require.Error(t, err) +} + +func TestEnsureReadDesire_RecreatesOnReadDesireGone(t *testing.T) { + ctx := context.Background() + store := &spyDeleteReadDesireStore{SpecStore: newMemoryStore()} + c := newTestClient(store) + + readID := desire.Identity{ + ManagementCluster: testManagementCluster, Type: desire.TypeRead, + Resource: testResource, Namespace: testNamespace, Name: testName, + } + + require.NoError(t, c.ensureReadDesire(ctx, readID, "v1", false)) + require.NoError(t, store.DeleteReadDesire(ctx, readID, c.owner, 1), "simulate an external deletion of the read desire") + require.NoError(t, c.ensureReadDesire(ctx, readID, "v1", false), + "an externally deleted read desire must be re-created on the next ensure") +} + +func TestApplyResource_VersionConflictSurfacesAsError(t *testing.T) { + ctx := context.Background() + memStore := newMemoryStore() + c := newTestClient(memStore) + + _, err := c.ApplyResource(ctx, configMapManifest(1), nil, testTransportContext()) + require.NoError(t, err) + + applyID := desire.Identity{ + ManagementCluster: testManagementCluster, Type: desire.TypeApply, + Resource: testResource, Namespace: testNamespace, Name: testName, + } + // Bump the real version out from under the (stale-reporting) view the + // client will see next, so its CAS write genuinely conflicts. + _, err = memStore.UpdateApplyDesireSpec( + ctx, applyID, desire.ApplySpec{KubeContent: configMapManifest(2)}, testOwner, 1) + require.NoError(t, err) + + staleClient := newTestClient(&staleApplyVersionStore{SpecStore: memStore}) + _, err = staleClient.ApplyResource(ctx, configMapManifest(3), nil, testTransportContext()) + require.Error(t, err) + assert.True(t, errors.Is(err, desire.ErrVersionConflict)) +} + +func TestEnsureReadDesire_RecreateAndVersionChangeMatrix(t *testing.T) { + tests := []struct { + name string + newVersion string + wantVersion string + wantDeleteCalls int + recreate bool + }{ + { + name: "recreate=false, same version: skipped", + recreate: false, + newVersion: "v1", + wantDeleteCalls: 0, + wantVersion: "v1", + }, + { + name: "recreate=false, version changed: skipped", + recreate: false, + newVersion: "v2", + wantDeleteCalls: 0, + wantVersion: "v1", + }, + { + name: "recreate=true, same version: skipped", + recreate: true, + newVersion: "v1", + wantDeleteCalls: 0, + wantVersion: "v1", + }, + { + name: "recreate=true, version changed: deletes and recreates", + recreate: true, + newVersion: "v2", + wantDeleteCalls: 1, + wantVersion: "v2", + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + ctx := context.Background() + store := &spyDeleteReadDesireStore{SpecStore: newMemoryStore()} + c := newTestClient(store) + + readID := desire.Identity{ + ManagementCluster: testManagementCluster, Type: desire.TypeRead, + Resource: testResource, Namespace: testNamespace, Name: testName, + } + + require.NoError(t, c.ensureReadDesire(ctx, readID, "v1", tt.recreate)) + require.NoError(t, c.ensureReadDesire(ctx, readID, tt.newVersion, tt.recreate)) + + assert.Equal(t, tt.wantDeleteCalls, store.deleteReadDesireCalls) + + rd, err := store.GetReadDesire(ctx, readID) + require.NoError(t, err) + assert.Equal(t, tt.wantVersion, rd.TargetVersion) + }) + } +} diff --git a/internal/desireclient/client.go b/internal/desireclient/client.go new file mode 100644 index 00000000..37356326 --- /dev/null +++ b/internal/desireclient/client.go @@ -0,0 +1,31 @@ +// Package desireclient implements transportclient.TransportClient against the +// desire-store contract from github.com/openshift-hyperfleet/hyperfleet-applier. +// It is the producer half of desire-based delivery: the adapter writes intent +// (apply/delete desires) and reads mirrored status (read desires) through the +// store, while a separate applier reconciles that intent against the target +// cluster. See docs/adapter-authoring-guide.md's "Desire transport" section +// for the contract this client implements. +package desireclient + +import ( + "github.com/openshift-hyperfleet/hyperfleet-adapter/internal/transportclient" + "github.com/openshift-hyperfleet/hyperfleet-adapter/pkg/logger" + "github.com/openshift-hyperfleet/hyperfleet-applier/pkg/desire" +) + +// Client implements transportclient.TransportClient against a desire store. +type Client struct { + store desire.SpecStore + log logger.Logger + owner string +} + +// NewClient builds a desire-backed transport client. store is the producer +// surface of the desire contract (never StatusStore, which is the applier's +// own status-writing surface). owner identifies this adapter as the writer +// for ownership/CAS checks on the store. +func NewClient(store desire.SpecStore, owner string, log logger.Logger) *Client { + return &Client{store: store, owner: owner, log: log} +} + +var _ transportclient.TransportClient = (*Client)(nil) diff --git a/internal/desireclient/delete.go b/internal/desireclient/delete.go new file mode 100644 index 00000000..1468801b --- /dev/null +++ b/internal/desireclient/delete.go @@ -0,0 +1,61 @@ +package desireclient + +import ( + "context" + "errors" + "fmt" + + "github.com/openshift-hyperfleet/hyperfleet-adapter/internal/transportclient" + "github.com/openshift-hyperfleet/hyperfleet-applier/pkg/desire" + "k8s.io/apimachinery/pkg/runtime/schema" +) + +// DeleteResource implements transportclient.TransportClient. It ensures a +// paired read desire exists, then posts a delete desire — CreateDeleteDesire +// atomically removes any sibling apply desire for the same target (see +// desire.SpecStore), so nothing re-applies. The read desire is created first +// so a failure there leaves the apply desire (if any) untouched — a safe, +// retryable state — rather than risking a delete desire with no read desire +// to observe its disappearance through, which a CreateDeleteDesire-then-read +// ordering could leave behind if the read-desire write failed afterward. +// +// opts is unused: propagation policy has no equivalent in the desire model +// (the applier owns deletion semantics against the target cluster), the same +// way maestroclient ignores it for ManifestWork deletes. +func (c *Client) DeleteResource( + ctx context.Context, + gvk schema.GroupVersionKind, + namespace, name string, + opts *transportclient.DeleteOptions, + target transportclient.TransportContext, +) error { + tc, err := resolveTransportContext(target) + if err != nil { + return err + } + + deleteID, err := buildIdentity(tc, desire.TypeDelete, gvk, namespace, name) + if err != nil { + return err + } + readID, err := buildIdentity(tc, desire.TypeRead, gvk, namespace, name) + if err != nil { + return err + } + + // recreate=false makes sure an existing readdesire is not recreated for a + // delete desire that will be processed soon. + if err = c.ensureReadDesire(ctx, readID, gvk.Version, false); err != nil { + return fmt.Errorf( + "desireclient: failed to create paired read desire for %s/%s: %w", gvk.Kind, name, err) + } + + if _, err = c.store.CreateDeleteDesire(ctx, desire.DeleteDesire{ + Identity: deleteID, + Owner: c.owner, + }); err != nil && !errors.Is(err, desire.ErrAlreadyExists) { + return fmt.Errorf("desireclient: failed to create delete desire for %s/%s: %w", namespace, name, err) + } + + return nil +} diff --git a/internal/desireclient/delete_test.go b/internal/desireclient/delete_test.go new file mode 100644 index 00000000..74dde324 --- /dev/null +++ b/internal/desireclient/delete_test.go @@ -0,0 +1,221 @@ +package desireclient + +import ( + "context" + "errors" + "testing" + + "github.com/openshift-hyperfleet/hyperfleet-applier/pkg/desire" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +func TestDeleteResource_RemovesApplyKeepsRead(t *testing.T) { + ctx := context.Background() + store := newMemoryStore() + c := newTestClient(store) + + _, err := c.ApplyResource(ctx, configMapManifest(1), nil, testTransportContext()) + require.NoError(t, err) + + applyID := desire.Identity{ + ManagementCluster: testManagementCluster, Type: desire.TypeApply, + Resource: testResource, Namespace: testNamespace, Name: testName, + } + readID := applyID + readID.Type = desire.TypeRead + deleteID := applyID + deleteID.Type = desire.TypeDelete + + err = c.DeleteResource(ctx, testGVK(), testNamespace, testName, nil, testTransportContext()) + require.NoError(t, err) + + _, err = store.GetApplyDesire(ctx, applyID) + assert.True(t, errors.Is(err, desire.ErrNotFound), "apply desire must be removed so nothing re-applies") + + _, err = store.GetDeleteDesire(ctx, deleteID) + require.NoError(t, err, "delete desire must be posted") + + _, err = store.GetReadDesire(ctx, readID) + require.NoError(t, err, "read desire must be left in place so disappearance stays observable") +} + +func TestDeleteResource_WithoutApply_PostsDelete(t *testing.T) { + ctx := context.Background() + store := newMemoryStore() + c := newTestClient(store) + + err := c.DeleteResource(ctx, testGVK(), testNamespace, testName, nil, testTransportContext()) + require.NoError(t, err) + + deleteID := desire.Identity{ + ManagementCluster: testManagementCluster, Type: desire.TypeDelete, + Resource: testResource, Namespace: testNamespace, Name: testName, + } + _, err = store.GetDeleteDesire(ctx, deleteID) + require.NoError(t, err) +} + +func TestDeleteResource_WithoutApply_PairsRead(t *testing.T) { + ctx := context.Background() + store := newMemoryStore() + c := newTestClient(store) + + // No ApplyResource call first: this identity was never applied through + // this client, so no read desire was ever paired in. + err := c.DeleteResource(ctx, testGVK(), testNamespace, testName, nil, testTransportContext()) + require.NoError(t, err) + + readID := desire.Identity{ + ManagementCluster: testManagementCluster, Type: desire.TypeRead, + Resource: testResource, Namespace: testNamespace, Name: testName, + } + readDesire, err := store.GetReadDesire(ctx, readID) + require.NoError(t, err, "paired read desire must be auto-created even without a prior apply") + assert.Equal(t, "v1", readDesire.TargetVersion) + assert.Equal(t, testOwner, readDesire.Owner) +} + +func TestDeleteResource_ReadFailure_ReturnsError(t *testing.T) { + ctx := context.Background() + store := &failingReadDesireStore{SpecStore: newMemoryStore()} + c := newTestClient(store) + + err := c.DeleteResource(ctx, testGVK(), testNamespace, testName, nil, testTransportContext()) + require.Error(t, err, "a read-desire pairing failure must surface as an error, not be swallowed") + assert.Contains(t, err.Error(), "paired read desire") +} + +// TestDeleteResource_ReadFailure_LeavesApplyIntact verifies +// the failure-injection contract for the read-before-delete ordering: an +// existing apply desire must survive a read-desire pairing failure untouched, +// and no delete desire must have been created, since the apply->delete +// transition (CreateDeleteDesire) never runs if the paired read desire can't +// be ensured first. +func TestDeleteResource_ReadFailure_LeavesApplyIntact(t *testing.T) { + ctx := context.Background() + inner := newMemoryStore() + + // Set up the pre-existing apply desire through a working client first — + // ApplyResource itself pairs a read desire via the same store call the + // failing wrapper below targets, so it must succeed here. + c := newTestClient(inner) + _, err := c.ApplyResource(ctx, configMapManifest(1), nil, testTransportContext()) + require.NoError(t, err) + + applyID := desire.Identity{ + ManagementCluster: testManagementCluster, Type: desire.TypeApply, + Resource: testResource, Namespace: testNamespace, Name: testName, + } + deleteID := applyID + deleteID.Type = desire.TypeDelete + + failingClient := newTestClient(&failingReadDesireStore{SpecStore: inner}) + err = failingClient.DeleteResource(ctx, testGVK(), testNamespace, testName, nil, testTransportContext()) + require.Error(t, err, "a read-desire pairing failure must surface as an error") + + _, err = inner.GetApplyDesire(ctx, applyID) + assert.NoError(t, err, + "apply desire must be untouched: the apply->delete transition must not run before the read desire is ensured") + + _, err = inner.GetDeleteDesire(ctx, deleteID) + assert.True(t, errors.Is(err, desire.ErrNotFound), + "delete desire must not exist: CreateDeleteDesire must not run when read-desire pairing fails first") +} + +// TestDeleteResource_DeleteFailure_LeavesApplyAndReadIntact +// verifies the second failure-injection path: if CreateDeleteDesire fails +// after the read desire has already been ensured, the apply desire must +// still be present (CreateDeleteDesire never got a chance to atomically +// remove it) and the read desire must remain — a safe, retryable state +// rather than an orphaned delete desire with no observable disappearance. +func TestDeleteResource_DeleteFailure_LeavesApplyAndReadIntact(t *testing.T) { + ctx := context.Background() + store := &failingCreateDeleteDesireStore{SpecStore: newMemoryStore()} + c := newTestClient(store) + + _, err := c.ApplyResource(ctx, configMapManifest(1), nil, testTransportContext()) + require.NoError(t, err) + + applyID := desire.Identity{ + ManagementCluster: testManagementCluster, Type: desire.TypeApply, + Resource: testResource, Namespace: testNamespace, Name: testName, + } + readID := applyID + readID.Type = desire.TypeRead + deleteID := applyID + deleteID.Type = desire.TypeDelete + + err = c.DeleteResource(ctx, testGVK(), testNamespace, testName, nil, testTransportContext()) + require.Error(t, err, "a delete-desire creation failure must surface as an error") + assert.Contains(t, err.Error(), "failed to create delete desire") + + _, err = store.GetApplyDesire(ctx, applyID) + assert.NoError(t, err, + "apply desire must still exist: CreateDeleteDesire failed before it could atomically remove it") + + _, err = store.GetReadDesire(ctx, readID) + assert.NoError(t, err, "read desire must still exist: it was ensured before the failed transition") + + _, err = store.GetDeleteDesire(ctx, deleteID) + assert.True(t, errors.Is(err, desire.ErrNotFound), "delete desire must not exist since its creation failed") +} + +func TestDeleteResource_DoesNotRecreateExistingReadDesire(t *testing.T) { + ctx := context.Background() + store := &spyDeleteReadDesireStore{SpecStore: newMemoryStore()} + c := newTestClient(store) + + readID := desire.Identity{ + ManagementCluster: testManagementCluster, Type: desire.TypeRead, + Resource: testResource, Namespace: testNamespace, Name: testName, + } + // Pre-existing read desire targets a different version than the GVK + // DeleteResource will be called with below. + require.NoError(t, c.ensureReadDesire(ctx, readID, "v1", false)) + + deleteGVK := testGVK() + deleteGVK.Version = "v1beta1" + + err := c.DeleteResource(ctx, deleteGVK, testNamespace, testName, nil, testTransportContext()) + require.NoError(t, err) + + assert.Zero(t, store.deleteReadDesireCalls, + "DeleteResource must never recreate the read desire, even on a version mismatch") + + rd, err := store.GetReadDesire(ctx, readID) + require.NoError(t, err) + assert.Equal(t, "v1", rd.TargetVersion, "read desire targetVersion must be left untouched by delete") +} + +func TestDeleteResource_RecreatesReadDesireIfExternallyDeleted(t *testing.T) { + ctx := context.Background() + store := newMemoryStore() + c := newTestClient(store) + + readID := desire.Identity{ + ManagementCluster: testManagementCluster, Type: desire.TypeRead, + Resource: testResource, Namespace: testNamespace, Name: testName, + } + require.NoError(t, c.ensureReadDesire(ctx, readID, "v1", false)) + + rd, err := store.GetReadDesire(ctx, readID) + require.NoError(t, err) + require.NoError(t, store.DeleteReadDesire(ctx, readID, c.owner, rd.Version), + "simulate an external deletion of the read desire") + + err = c.DeleteResource(ctx, testGVK(), testNamespace, testName, nil, testTransportContext()) + require.NoError(t, err) + + _, err = store.GetReadDesire(ctx, readID) + require.NoError(t, err, + "an externally deleted read desire must be re-created by delete, even though recreate=false") +} + +func TestDeleteResource_RequiresTransportContext(t *testing.T) { + ctx := context.Background() + c := newTestClient(newMemoryStore()) + + err := c.DeleteResource(ctx, testGVK(), testNamespace, testName, nil, nil) + require.Error(t, err) +} diff --git a/internal/desireclient/desireclient_test.go b/internal/desireclient/desireclient_test.go new file mode 100644 index 00000000..9bf1f56c --- /dev/null +++ b/internal/desireclient/desireclient_test.go @@ -0,0 +1,78 @@ +package desireclient + +import ( + "context" + "errors" + "fmt" + + "github.com/openshift-hyperfleet/hyperfleet-adapter/pkg/constants" + "github.com/openshift-hyperfleet/hyperfleet-adapter/pkg/logger" + "github.com/openshift-hyperfleet/hyperfleet-applier/pkg/desire" + "github.com/openshift-hyperfleet/hyperfleet-applier/pkg/desire/store/memory" + "k8s.io/apimachinery/pkg/runtime/schema" +) + +const ( + testManagementCluster = "mgmt-cluster-01" + testResource = "configmaps" + testNamespace = "default" + testName = "my-config" + testOwner = "hyperfleet-adapter" +) + +func newTestClient(store desire.SpecStore) *Client { + return NewClient(store, testOwner, logger.NewTestLogger()) +} + +func testTransportContext() *TransportContext { + return &TransportContext{ManagementCluster: testManagementCluster, Resource: testResource} +} + +func testGVK() schema.GroupVersionKind { + return schema.GroupVersionKind{Version: "v1", Kind: "ConfigMap"} +} + +// configMapManifest builds a minimal ConfigMap manifest carrying the +// hyperfleet.io/generation annotation task authors are required to set +// (docs/adapter-authoring-guide.md:498). +func configMapManifest(generation int64) []byte { + return fmt.Appendf(nil, `{ + "apiVersion": "v1", + "kind": "ConfigMap", + "metadata": { + "name": %q, + "namespace": %q, + "annotations": {%q: %q} + }, + "data": {"key": "value"} + }`, testName, testNamespace, constants.AnnotationGeneration, fmt.Sprint(generation)) +} + +// failingReadDesireStore wraps a real SpecStore but forces CreateReadDesire +// to fail, simulating a read-desire pairing failure. +type failingReadDesireStore struct { + desire.SpecStore +} + +func (f *failingReadDesireStore) CreateReadDesire( + _ context.Context, _ desire.ReadDesire, +) (desire.ReadDesire, error) { + return desire.ReadDesire{}, errors.New("boom: read desire store unavailable") +} + +// failingCreateDeleteDesireStore wraps a real SpecStore but forces +// CreateDeleteDesire to fail, simulating a failure in the apply-to-delete +// transition after the paired read desire has already been ensured. +type failingCreateDeleteDesireStore struct { + desire.SpecStore +} + +func (f *failingCreateDeleteDesireStore) CreateDeleteDesire( + _ context.Context, _ desire.DeleteDesire, +) (desire.DeleteDesire, error) { + return desire.DeleteDesire{}, errors.New("boom: delete desire store unavailable") +} + +func newMemoryStore() *memory.Store { + return memory.New() +} diff --git a/internal/desireclient/discover.go b/internal/desireclient/discover.go new file mode 100644 index 00000000..c55fe61d --- /dev/null +++ b/internal/desireclient/discover.go @@ -0,0 +1,73 @@ +package desireclient + +import ( + "context" + "errors" + "fmt" + + "github.com/openshift-hyperfleet/hyperfleet-adapter/internal/manifest" + "github.com/openshift-hyperfleet/hyperfleet-adapter/internal/transportclient" + "github.com/openshift-hyperfleet/hyperfleet-adapter/pkg/logger" + apierrors "k8s.io/apimachinery/pkg/api/errors" + "k8s.io/apimachinery/pkg/apis/meta/v1/unstructured" + "k8s.io/apimachinery/pkg/runtime/schema" +) + +// DiscoverResources implements transportclient.TransportClient. It lists +// every read desire in the target partition and filters by GVK and the +// discovery criteria. desire.Identity carries no labels, so label-selector +// discovery must scan client-side, the same shape as +// maestroclient.DiscoverResources scanning ManifestWorks. +func (c *Client) DiscoverResources( + ctx context.Context, + gvk schema.GroupVersionKind, + discovery manifest.Discovery, + target transportclient.TransportContext, +) (*unstructured.UnstructuredList, error) { + tc, err := resolveTransportContext(target) + if err != nil { + return nil, err + } + + reads, err := c.store.ListReadDesires(ctx, tc.ManagementCluster) + if err != nil { + return nil, fmt.Errorf("desireclient: failed to list read desires for partition %q: %w", tc.ManagementCluster, err) + } + + list := &unstructured.UnstructuredList{} + for _, rd := range reads { + if rd.Identity.Group != gvk.Group || rd.Identity.Resource != tc.Resource { + continue + } + + // Route through the same three-way interpretation GetResource uses, so a + // single item's outcome here can never drift from what a Get on that same + // identity would report. + obj, err := c.decodeReadDesire(gvk, rd.Identity.Namespace, rd.Identity.Name, rd) + switch { + case errors.Is(err, ErrNotSyncedYet): + // Not yet synced — nothing to match discovery criteria against, + // same as a live List not yet showing a slow-to-create resource. + continue + case apierrors.IsNotFound(err): + // Confirmed absent — same as a live List not showing a deleted resource. + continue + case err != nil: + // Undecodable KubeContent (whether from a synced True condition or + // a retained mirror on a KubeAPIError/PreCheckFailed False + // condition) is a store invariant violation for that one record, + // not legitimate transience — but it must not fail discovery for + // every other resource in the partition. + errCtx := logger.WithErrorField(ctx, err) + c.log.Errorf(errCtx, "desireclient: discovery skipping %s/%s: %v", + rd.Identity.Namespace, rd.Identity.Name, err) + continue + } + + if manifest.MatchesDiscoveryCriteria(obj, discovery) { + list.Items = append(list.Items, *obj) + } + } + + return list, nil +} diff --git a/internal/desireclient/discover_test.go b/internal/desireclient/discover_test.go new file mode 100644 index 00000000..f5911c68 --- /dev/null +++ b/internal/desireclient/discover_test.go @@ -0,0 +1,226 @@ +package desireclient + +import ( + "context" + "testing" + + "github.com/openshift-hyperfleet/hyperfleet-adapter/internal/manifest" + "github.com/openshift-hyperfleet/hyperfleet-applier/pkg/desire" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" +) + +func TestDiscoverResources_EmptyPartitionReturnsEmptyList(t *testing.T) { + ctx := context.Background() + c := newTestClient(newMemoryStore()) + + list, err := c.DiscoverResources(ctx, testGVK(), &manifest.DiscoveryConfig{}, testTransportContext()) + require.NoError(t, err) + assert.Empty(t, list.Items) +} + +func TestDiscoverResources_ReturnsSyncedResourceByName(t *testing.T) { + ctx := context.Background() + store := newMemoryStore() + c := newTestClient(store) + + putSyncedReadDesire(t, ctx, store, testNamespace, testName, configMapManifest(1)) + + list, err := c.DiscoverResources(ctx, testGVK(), &manifest.DiscoveryConfig{ByName: testName}, testTransportContext()) + require.NoError(t, err) + require.Len(t, list.Items, 1) + assert.Equal(t, testName, list.Items[0].GetName()) +} + +func TestDiscoverResources_ByNameExcludesNonMatchingName(t *testing.T) { + ctx := context.Background() + store := newMemoryStore() + c := newTestClient(store) + + putSyncedReadDesire(t, ctx, store, testNamespace, testName, configMapManifest(1)) + + discovery := &manifest.DiscoveryConfig{ByName: "other-name"} + list, err := c.DiscoverResources(ctx, testGVK(), discovery, testTransportContext()) + require.NoError(t, err) + assert.Empty(t, list.Items, "discovery criteria must filter out non-matching names") +} + +func TestDiscoverResources_LabelSelectorMatchesSubset(t *testing.T) { + ctx := context.Background() + store := newMemoryStore() + c := newTestClient(store) + + labeledManifest := []byte(`{ + "apiVersion": "v1", "kind": "ConfigMap", + "metadata": {"name": "labeled", "namespace": "default", "labels": {"app": "myapp"}} + }`) + unlabeledManifest := []byte(`{ + "apiVersion": "v1", "kind": "ConfigMap", + "metadata": {"name": "unlabeled", "namespace": "default"} + }`) + putSyncedReadDesire(t, ctx, store, "labeled", "labeled", labeledManifest) + putSyncedReadDesire(t, ctx, store, "unlabeled", "unlabeled", unlabeledManifest) + + list, err := c.DiscoverResources(ctx, testGVK(), + &manifest.DiscoveryConfig{LabelSelector: "app=myapp"}, testTransportContext()) + require.NoError(t, err) + require.Len(t, list.Items, 1) + assert.Equal(t, "labeled", list.Items[0].GetName()) +} + +func TestDiscoverResources_SkipsNotYetSyncedResource(t *testing.T) { + ctx := context.Background() + store := newMemoryStore() + c := newTestClient(store) + + id := desire.Identity{ + ManagementCluster: testManagementCluster, Type: desire.TypeRead, + Resource: testResource, Namespace: testNamespace, Name: testName, + } + _, err := store.CreateReadDesire(ctx, desire.ReadDesire{Identity: id, Owner: testOwner, TargetVersion: "v1"}) + require.NoError(t, err) + + list, err := c.DiscoverResources(ctx, testGVK(), &manifest.DiscoveryConfig{}, testTransportContext()) + require.NoError(t, err) + assert.Empty(t, list.Items, "a read desire with no Successful condition yet must not appear in discovery") +} + +func TestDiscoverResources_SurfacesRetainedMirrorOnFailedRead(t *testing.T) { + ctx := context.Background() + store := newMemoryStore() + c := newTestClient(store) + + putKubeAPIErrorReadDesire(t, ctx, store, testNamespace, testName, configMapManifest(1)) + + list, err := c.DiscoverResources(ctx, testGVK(), &manifest.DiscoveryConfig{}, testTransportContext()) + require.NoError(t, err) + require.Len(t, list.Items, 1, + "Successful=False/KubeAPIError retains the last mirrored content - it is still 'present', just possibly stale") + assert.Equal(t, testName, list.Items[0].GetName()) +} + +func TestDiscoverResources_SkipsFailedReadWithNoRetainedMirror(t *testing.T) { + ctx := context.Background() + store := newMemoryStore() + c := newTestClient(store) + + putKubeAPIErrorReadDesire(t, ctx, store, testNamespace, testName, nil) + + list, err := c.DiscoverResources(ctx, testGVK(), &manifest.DiscoveryConfig{}, testTransportContext()) + require.NoError(t, err) + assert.Empty(t, list.Items, "no content has ever been mirrored, so there is nothing to surface") +} + +func TestDiscoverResources_SkipsNotFoundFalseDesire(t *testing.T) { + ctx := context.Background() + store := newMemoryStore() + c := newTestClient(store) + + putNotFoundReadDesire(t, ctx, store, testNamespace, testName) + + list, err := c.DiscoverResources(ctx, testGVK(), &manifest.DiscoveryConfig{}, testTransportContext()) + require.NoError(t, err) + assert.Empty(t, list.Items, "Successful=False/NotFound desire must not appear in discovery") +} + +func TestDiscoverResources_SkipsUndecodableContentButKeepsOthers(t *testing.T) { + ctx := context.Background() + store := newMemoryStore() + c := newTestClient(store) + + putSyncedReadDesire(t, ctx, store, testNamespace, "bad", []byte("not-json")) + putSyncedReadDesire(t, ctx, store, testNamespace, "good", configMapManifest(1)) + + list, err := c.DiscoverResources(ctx, testGVK(), &manifest.DiscoveryConfig{}, testTransportContext()) + require.NoError(t, err, "a single bad record must not fail discovery for the whole partition") + require.Len(t, list.Items, 1) + assert.Equal(t, testName, list.Items[0].GetName()) +} + +func TestDiscoverResources_FiltersOutOtherResourceType(t *testing.T) { + ctx := context.Background() + store := newMemoryStore() + c := newTestClient(store) + + putSyncedReadDesire(t, ctx, store, testNamespace, testName, configMapManifest(1)) + + otherContext := &TransportContext{ManagementCluster: testManagementCluster, Resource: "secrets"} + list, err := c.DiscoverResources(ctx, testGVK(), &manifest.DiscoveryConfig{}, otherContext) + require.NoError(t, err) + assert.Empty(t, list.Items, "a read desire for a different plural resource must not appear") +} + +func TestDiscoverResources_FiltersOutOtherGroup(t *testing.T) { + ctx := context.Background() + store := newMemoryStore() + c := newTestClient(store) + + id := desire.Identity{ + ManagementCluster: testManagementCluster, Type: desire.TypeRead, + Group: "apps", Resource: testResource, Namespace: testNamespace, Name: testName, + } + _, err := store.CreateReadDesire(ctx, desire.ReadDesire{Identity: id, Owner: testOwner, TargetVersion: "v1"}) + require.NoError(t, err) + _, err = store.UpdateReadDesireStatus(ctx, id, desire.ReadStatus{ + Status: desire.Status{Conditions: []metav1.Condition{{ + Type: desire.TypeSuccessful, Status: metav1.ConditionTrue, Reason: desire.ReasonSynced, + }}}, + KubeContent: configMapManifest(1), + }) + require.NoError(t, err) + + // testGVK() has an empty Group ("core"), so a read desire recorded under + // group "apps" must not match even though Resource matches. + list, err := c.DiscoverResources(ctx, testGVK(), &manifest.DiscoveryConfig{}, testTransportContext()) + require.NoError(t, err) + assert.Empty(t, list.Items, "a read desire in a different API group must not appear") +} + +func TestDiscoverResources_ScopedToPartition(t *testing.T) { + ctx := context.Background() + store := newMemoryStore() + c := newTestClient(store) + + putSyncedReadDesire(t, ctx, store, testNamespace, testName, configMapManifest(1)) + + otherPartition := &TransportContext{ManagementCluster: "other-cluster", Resource: testResource} + list, err := c.DiscoverResources(ctx, testGVK(), &manifest.DiscoveryConfig{}, otherPartition) + require.NoError(t, err) + assert.Empty(t, list.Items, "discovery must not leak resources across management-cluster partitions") +} + +func TestDiscoverResources_MultipleMatches(t *testing.T) { + ctx := context.Background() + store := newMemoryStore() + c := newTestClient(store) + + first := []byte(`{"apiVersion": "v1", "kind": "ConfigMap", "metadata": {"name": "first", "namespace": "default"}}`) + second := []byte(`{"apiVersion": "v1", "kind": "ConfigMap", "metadata": {"name": "second", "namespace": "default"}}`) + putSyncedReadDesire(t, ctx, store, testNamespace, "first", first) + putSyncedReadDesire(t, ctx, store, testNamespace, "second", second) + + list, err := c.DiscoverResources(ctx, testGVK(), &manifest.DiscoveryConfig{}, testTransportContext()) + require.NoError(t, err) + require.Len(t, list.Items, 2) + names := []string{list.Items[0].GetName(), list.Items[1].GetName()} + assert.ElementsMatch(t, []string{"first", "second"}, names) +} + +func TestDiscoverResources_ListFailure_ReturnsError(t *testing.T) { + ctx := context.Background() + store := &failingListReadDesiresStore{SpecStore: newMemoryStore()} + c := newTestClient(store) + + _, err := c.DiscoverResources(ctx, testGVK(), &manifest.DiscoveryConfig{}, testTransportContext()) + require.Error(t, err) + assert.Contains(t, err.Error(), "failed to list read desires") +} + +func TestDiscoverResources_RequiresTransportContext(t *testing.T) { + ctx := context.Background() + c := newTestClient(newMemoryStore()) + + _, err := c.DiscoverResources(ctx, testGVK(), &manifest.DiscoveryConfig{}, nil) + require.Error(t, err) +} diff --git a/internal/desireclient/get.go b/internal/desireclient/get.go new file mode 100644 index 00000000..f281cea3 --- /dev/null +++ b/internal/desireclient/get.go @@ -0,0 +1,92 @@ +package desireclient + +import ( + "context" + "encoding/json" + "errors" + "fmt" + + "github.com/openshift-hyperfleet/hyperfleet-adapter/internal/transportclient" + "github.com/openshift-hyperfleet/hyperfleet-applier/pkg/desire" + apierrors "k8s.io/apimachinery/pkg/api/errors" + apimeta "k8s.io/apimachinery/pkg/api/meta" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/apis/meta/v1/unstructured" + "k8s.io/apimachinery/pkg/runtime/schema" +) + +// GetResource implements transportclient.TransportClient. It reads the +// mirrored live object from the read desire's status, distinguishing three +// outcomes per the eventual-consistency contract +// (docs/adapter-authoring-guide.md): not-synced-yet (ErrNotSyncedYet), +// confirmed-absent (apierrors.NewNotFound), and present (the mirrored +// object). +func (c *Client) GetResource( + ctx context.Context, + gvk schema.GroupVersionKind, + namespace, name string, + target transportclient.TransportContext, +) (*unstructured.Unstructured, error) { + tc, err := resolveTransportContext(target) + if err != nil { + return nil, err + } + + id, err := buildIdentity(tc, desire.TypeRead, gvk, namespace, name) + if err != nil { + return nil, err + } + + rd, err := c.store.GetReadDesire(ctx, id) + if errors.Is(err, desire.ErrNotFound) { + return nil, ErrNotSyncedYet + } + if err != nil { + return nil, fmt.Errorf("desireclient: failed to get read desire for %s/%s: %w", namespace, name, err) + } + + return c.decodeReadDesire(gvk, namespace, name, rd) +} + +// decodeReadDesire translates a ReadDesire's status conditions into the +// three-way outcome the eventual-consistency contract defines +// (docs/adapter-authoring-guide.md, "The eventual-consistency contract for +// remote reads"): no condition yet is not-synced-yet, Reason=NotFound is a +// confirmed absence, and everything else - including a transient failure +// (KubeAPIError/PreCheckFailed) - decodes whatever KubeContent the applier +// last mirrored. readdesire's status.go (applier) retains the prior mirror +// across a transient failure rather than clearing it, so that content is +// still "present" per the contract, just possibly stale; staleness is the +// caller's concern via the generation annotation, not this layer's. +func (c *Client) decodeReadDesire( + gvk schema.GroupVersionKind, namespace, name string, rd desire.ReadDesire, +) (*unstructured.Unstructured, error) { + cond := apimeta.FindStatusCondition(rd.Status.Conditions, desire.TypeSuccessful) + + switch { + case cond == nil: + // Read desire exists but the applier hasn't observed it yet. + return nil, ErrNotSyncedYet + + case cond.Status == metav1.ConditionFalse && cond.Reason == desire.ReasonNotFound: + return nil, apierrors.NewNotFound(schema.GroupResource{Group: gvk.Group, Resource: rd.Identity.Resource}, name) + + default: + return decodeKubeContent(rd.Status.KubeContent, gvk, rd.Identity.Resource, namespace, name) + } +} + +// decodeKubeContent decodes a synced read desire's mirrored content, or +// reports confirmed-absence when content is empty. +func decodeKubeContent( + kubeContent []byte, gvk schema.GroupVersionKind, resource, namespace, name string, +) (*unstructured.Unstructured, error) { + if len(kubeContent) == 0 { + return nil, apierrors.NewNotFound(schema.GroupResource{Group: gvk.Group, Resource: resource}, name) + } + obj := &unstructured.Unstructured{} + if err := json.Unmarshal(kubeContent, obj); err != nil { + return nil, fmt.Errorf("desireclient: failed to decode mirrored content for %s/%s: %w", namespace, name, err) + } + return obj, nil +} diff --git a/internal/desireclient/get_test.go b/internal/desireclient/get_test.go new file mode 100644 index 00000000..f049040d --- /dev/null +++ b/internal/desireclient/get_test.go @@ -0,0 +1,212 @@ +package desireclient + +import ( + "context" + "errors" + "testing" + + "github.com/openshift-hyperfleet/hyperfleet-applier/pkg/desire" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + apierrors "k8s.io/apimachinery/pkg/api/errors" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" +) + +func readIdentity() desire.Identity { + return desire.Identity{ + ManagementCluster: testManagementCluster, Type: desire.TypeRead, + Resource: testResource, Namespace: testNamespace, Name: testName, + } +} + +func TestGetResource_ReadDesireNotFoundIsNotSyncedYet(t *testing.T) { + ctx := context.Background() + c := newTestClient(newMemoryStore()) + + _, err := c.GetResource(ctx, testGVK(), testNamespace, testName, testTransportContext()) + require.Error(t, err) + assert.True(t, errors.Is(err, ErrNotSyncedYet)) + assert.False(t, apierrors.IsNotFound(err), "absent mirror must not collapse into NotFound") +} + +func TestGetResource_ReadDesireExistsNoResourceObservedIsNotSyncedYet(t *testing.T) { + ctx := context.Background() + store := newMemoryStore() + c := newTestClient(store) + + _, err := store.CreateReadDesire(ctx, desire.ReadDesire{ + Identity: readIdentity(), Owner: testOwner, TargetVersion: "v1", + }) + require.NoError(t, err) + + _, err = c.GetResource(ctx, testGVK(), testNamespace, testName, testTransportContext()) + require.Error(t, err) + assert.True(t, errors.Is(err, ErrNotSyncedYet), "no Successful condition yet must read as not-synced-yet") +} + +func TestGetResource_SyncedReturnsMirroredObject(t *testing.T) { + ctx := context.Background() + store := newMemoryStore() + c := newTestClient(store) + + putSyncedReadDesire(t, ctx, store, testNamespace, testName, configMapManifest(1)) + + obj, err := c.GetResource(ctx, testGVK(), testNamespace, testName, testTransportContext()) + require.NoError(t, err) + assert.Equal(t, testName, obj.GetName()) + assert.Equal(t, testNamespace, obj.GetNamespace()) +} + +func TestGetResource_ConfirmedNotFound(t *testing.T) { + ctx := context.Background() + store := newMemoryStore() + c := newTestClient(store) + + putConfirmedAbsentReadDesire(t, ctx, store, testNamespace, testName) + + _, err := c.GetResource(ctx, testGVK(), testNamespace, testName, testTransportContext()) + require.Error(t, err) + assert.True(t, apierrors.IsNotFound(err), "confirmed-gone must be a real NotFound, not ErrNotSyncedYet") + assert.False(t, errors.Is(err, ErrNotSyncedYet)) +} + +func TestGetResource_InvalidReadDesire(t *testing.T) { + ctx := context.Background() + store := newMemoryStore() + c := newTestClient(store) + + putInvalidReadDesire(t, ctx, store, testNamespace, testName) + + _, err := c.GetResource(ctx, testGVK(), testNamespace, testName, testTransportContext()) + require.Error(t, err) + assert.True(t, apierrors.IsNotFound(err), "confirmed-gone must be a real NotFound, not ErrNotSyncedYet") + assert.False(t, errors.Is(err, ErrNotSyncedYet)) +} + +func TestGetResource_K8sAPIErrorWithRetainedMirrorReturnsStaleContent(t *testing.T) { + ctx := context.Background() + store := newMemoryStore() + c := newTestClient(store) + + putKubeAPIErrorReadDesire(t, ctx, store, testNamespace, testName, configMapManifest(1)) + + obj, err := c.GetResource(ctx, testGVK(), testNamespace, testName, testTransportContext()) + require.NoError(t, err, + "readdesire.kubeAPIError retains the last mirrored content across a transient failure - "+ + "it is still 'present' per the eventual-consistency contract, just possibly stale") + assert.Equal(t, testName, obj.GetName()) +} + +func TestGetResource_K8sAPIErrorWithNoMirrorYetIsConfirmedAbsent(t *testing.T) { + ctx := context.Background() + store := newMemoryStore() + c := newTestClient(store) + + putKubeAPIErrorReadDesire(t, ctx, store, testNamespace, testName, nil) + + _, err := c.GetResource(ctx, testGVK(), testNamespace, testName, testTransportContext()) + require.Error(t, err) + assert.True(t, apierrors.IsNotFound(err), + "a transient failure with no content ever mirrored decodes as empty, "+ + "same as decodeKubeContent's own empty-content rule") +} + +func TestDecodeKubeContent(t *testing.T) { + tests := []struct { + name string + content []byte + wantNotFound bool + wantErr bool + }{ + {name: "empty content is confirmed absent", content: nil, wantNotFound: true}, + {name: "valid content decodes", content: configMapManifest(1)}, + {name: "invalid json is a decode error", content: []byte("not-json"), wantErr: true}, + } + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + obj, err := decodeKubeContent(tt.content, testGVK(), testResource, testNamespace, testName) + switch { + case tt.wantNotFound: + require.Error(t, err) + assert.True(t, apierrors.IsNotFound(err)) + case tt.wantErr: + require.Error(t, err) + assert.False(t, apierrors.IsNotFound(err), "a decode failure is not the same outcome as confirmed-absence") + default: + require.NoError(t, err) + assert.Equal(t, testName, obj.GetName()) + assert.Equal(t, testNamespace, obj.GetNamespace()) + } + }) + } +} + +func TestDecodeReadDesire(t *testing.T) { + tests := []struct { + name string + condition *metav1.Condition + content []byte + wantNotSynced bool + wantNotFound bool + wantContent bool + }{ + { + name: "no condition yet is not synced", + condition: nil, + wantNotSynced: true, + }, + { + name: "successful true decodes content", + condition: successfulCondition(metav1.ConditionTrue, desire.ReasonSynced), + content: configMapManifest(1), + wantContent: true, + }, + { + name: "successful true with empty content is confirmed absent", + condition: successfulCondition(metav1.ConditionTrue, desire.ReasonNotFound), + wantNotFound: true, + }, + { + name: "false with notfound reason is confirmed absent", + condition: successfulCondition(metav1.ConditionFalse, desire.ReasonNotFound), + wantNotFound: true, + }, + { + name: "false with other reason decodes the retained mirror when present", + condition: successfulCondition(metav1.ConditionFalse, desire.ReasonKubeAPIError), + content: configMapManifest(1), + wantContent: true, + }, + { + name: "false with other reason and no retained mirror is confirmed absent", + condition: successfulCondition(metav1.ConditionFalse, desire.ReasonKubeAPIError), + wantNotFound: true, + }, + } + + c := newTestClient(newMemoryStore()) + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + rd := desire.ReadDesire{Identity: readIdentity()} + if tt.condition != nil { + rd.Status = desire.ReadStatus{ + Status: desire.Status{Conditions: []metav1.Condition{*tt.condition}}, + KubeContent: tt.content, + } + } + + obj, err := c.decodeReadDesire(testGVK(), testNamespace, testName, rd) + switch { + case tt.wantNotSynced: + require.Error(t, err) + assert.True(t, errors.Is(err, ErrNotSyncedYet)) + case tt.wantNotFound: + require.Error(t, err) + assert.True(t, apierrors.IsNotFound(err)) + case tt.wantContent: + require.NoError(t, err) + assert.Equal(t, testName, obj.GetName()) + } + }) + } +} diff --git a/internal/desireclient/helpers_test.go b/internal/desireclient/helpers_test.go new file mode 100644 index 00000000..e5768b13 --- /dev/null +++ b/internal/desireclient/helpers_test.go @@ -0,0 +1,129 @@ +package desireclient + +import ( + "context" + "errors" + "testing" + + "github.com/openshift-hyperfleet/hyperfleet-applier/pkg/desire" + "github.com/openshift-hyperfleet/hyperfleet-applier/pkg/desire/store/memory" + "github.com/stretchr/testify/require" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" +) + +// putReadDesire creates a read desire with the given status. +func putReadDesire( + t *testing.T, ctx context.Context, store *memory.Store, namespace, name string, status desire.ReadStatus, +) { + t.Helper() + id := desire.Identity{ + ManagementCluster: testManagementCluster, Type: desire.TypeRead, + Resource: testResource, Namespace: namespace, Name: name, + } + _, err := store.CreateReadDesire(ctx, desire.ReadDesire{Identity: id, Owner: testOwner, TargetVersion: "v1"}) + require.NoError(t, err) + _, err = store.UpdateReadDesireStatus(ctx, id, status) + require.NoError(t, err) +} + +// putConfirmedAbsentReadDesire creates a read desire marked as not found. +func putConfirmedAbsentReadDesire(t *testing.T, ctx context.Context, store *memory.Store, namespace, name string) { + t.Helper() + putReadDesire(t, ctx, store, namespace, name, desire.ReadStatus{ + Status: desire.Status{Conditions: []metav1.Condition{{ + Type: desire.TypeSuccessful, Status: metav1.ConditionFalse, Reason: desire.ReasonNotFound, + }}}, + }) +} + +// putSyncedReadDesire creates a read desire with synced content. +func putSyncedReadDesire( + t *testing.T, ctx context.Context, store *memory.Store, namespace, name string, content []byte, +) { + t.Helper() + putReadDesire(t, ctx, store, namespace, name, desire.ReadStatus{ + Status: desire.Status{Conditions: []metav1.Condition{{ + Type: desire.TypeSuccessful, Status: metav1.ConditionTrue, Reason: desire.ReasonSynced, + }}}, + KubeContent: content, + }) +} + +// putNotFoundReadDesire creates a read desire marked as not found. +func putNotFoundReadDesire(t *testing.T, ctx context.Context, store *memory.Store, namespace, name string) { + t.Helper() + putReadDesire(t, ctx, store, namespace, name, desire.ReadStatus{ + Status: desire.Status{Conditions: []metav1.Condition{{ + Type: desire.TypeSuccessful, Status: metav1.ConditionFalse, Reason: desire.ReasonNotFound, + }}}, + }) +} + +// putInvalidReadDesire creates a read desire with invalid content (successful=true but notfound reason). +func putInvalidReadDesire(t *testing.T, ctx context.Context, store *memory.Store, namespace, name string) { + t.Helper() + putReadDesire(t, ctx, store, namespace, name, desire.ReadStatus{ + Status: desire.Status{Conditions: []metav1.Condition{{ + Type: desire.TypeSuccessful, Status: metav1.ConditionTrue, Reason: desire.ReasonNotFound, + }}}, + }) +} + +// putKubeAPIErrorReadDesire creates a read desire with a transient kube API error. +func putKubeAPIErrorReadDesire( + t *testing.T, ctx context.Context, store *memory.Store, namespace, name string, content []byte, +) { + t.Helper() + putReadDesire(t, ctx, store, namespace, name, desire.ReadStatus{ + Status: desire.Status{Conditions: []metav1.Condition{{ + Type: desire.TypeSuccessful, Status: metav1.ConditionFalse, Reason: desire.ReasonKubeAPIError, + }}}, + KubeContent: content, + }) +} + +// successfulCondition builds the single summary condition every desire carries. +func successfulCondition(status metav1.ConditionStatus, reason string) *metav1.Condition { + return &metav1.Condition{Type: desire.TypeSuccessful, Status: status, Reason: reason} +} + +// failingListReadDesiresStore wraps a real SpecStore but forces +// ListReadDesires to fail, simulating a store-level outage during discovery. +type failingListReadDesiresStore struct { + desire.SpecStore +} + +func (f *failingListReadDesiresStore) ListReadDesires( + _ context.Context, _ string, +) ([]desire.ReadDesire, error) { + return nil, errors.New("boom: read desire store unavailable") +} + +// spyDeleteReadDesireStore counts DeleteReadDesire calls so tests can assert +// whether ensureReadDesire actually attempted a recreate. +type spyDeleteReadDesireStore struct { + desire.SpecStore + deleteReadDesireCalls int +} + +func (s *spyDeleteReadDesireStore) DeleteReadDesire( + ctx context.Context, id desire.Identity, owner string, version int64, +) error { + s.deleteReadDesireCalls++ + return s.SpecStore.DeleteReadDesire(ctx, id, owner, version) +} + +// staleApplyVersionStore wraps a real SpecStore but returns a stale version +// on GetApplyDesire to simulate the case where an external client has +// concurrently updated the apply desire while this one is computing. +type staleApplyVersionStore struct { + desire.SpecStore +} + +func (s *staleApplyVersionStore) GetApplyDesire(ctx context.Context, id desire.Identity) (desire.ApplyDesire, error) { + ad, err := s.SpecStore.GetApplyDesire(ctx, id) + if err == nil { + ad.Version++ + } + return ad, err +} diff --git a/internal/desireclient/types.go b/internal/desireclient/types.go new file mode 100644 index 00000000..e8f9bfe1 --- /dev/null +++ b/internal/desireclient/types.go @@ -0,0 +1,104 @@ +package desireclient + +import ( + "encoding/json" + "errors" + "fmt" + + "github.com/openshift-hyperfleet/hyperfleet-adapter/internal/manifest" + "github.com/openshift-hyperfleet/hyperfleet-adapter/internal/transportclient" + "github.com/openshift-hyperfleet/hyperfleet-applier/pkg/desire" + "k8s.io/apimachinery/pkg/apis/meta/v1/unstructured" + "k8s.io/apimachinery/pkg/runtime/schema" + "sigs.k8s.io/yaml" +) + +// TransportContext carries per-request routing information for the desire +// transport backend. Pass this as the TransportContext (any) in ApplyResource, +// GetResource, DiscoverResources, and DeleteResource. +type TransportContext struct { + // ManagementCluster is the target managed cluster identifier — the desire + // store's partition key. Required for all operations. + ManagementCluster string + // Resource is the plural Kubernetes resource type (e.g. "networkpolicies"), + // declared in the task config alongside the manifest. The desire store's + // Identity is keyed by this plural form directly (see + // pkg/desire.Identity.Resource in hyperfleet-applier); the adapter never + // derives it from the manifest's Kind — no RESTMapper is used or needed. + // Required for all operations. + Resource string +} + +// ErrNotSyncedYet indicates the read desire's mirror has not yet been +// populated by the applier — a transient, non-terminal outcome distinct from +// a confirmed-absent resource (reported via apierrors.NewNotFound). Per the +// eventual-consistency contract (docs/adapter-authoring-guide.md), callers +// should treat this as "not converged yet" rather than a failure. Use +// errors.Is to check for it. +var ErrNotSyncedYet = errors.New("desireclient: resource not synced yet") + +// resolveTransportContext type-asserts the generic TransportContext and +// validates both fields are set. +func resolveTransportContext(target transportclient.TransportContext) (*TransportContext, error) { + tc, ok := target.(*TransportContext) + if !ok || tc == nil { + return nil, fmt.Errorf("desireclient: TransportContext with ManagementCluster and Resource is required") + } + if tc.ManagementCluster == "" { + return nil, fmt.Errorf("desireclient: TransportContext.ManagementCluster is required") + } + if tc.Resource == "" { + return nil, fmt.Errorf("desireclient: TransportContext.Resource is required") + } + return tc, nil +} + +// buildIdentity is the single shared helper for assembling a desire.Identity +// from a resolved TransportContext, GVK, namespace, and name. All four +// TransportClient methods route through this rather than duplicating identity +// construction. +func buildIdentity( + tc *TransportContext, dtype desire.DesireType, gvk schema.GroupVersionKind, namespace, name string, +) (desire.Identity, error) { + id := desire.Identity{ + ManagementCluster: tc.ManagementCluster, + Type: dtype, + Group: gvk.Group, + Resource: tc.Resource, + Namespace: namespace, + Name: name, + } + if err := id.Validate(); err != nil { + return desire.Identity{}, fmt.Errorf("desireclient: invalid identity: %w", err) + } + return id, nil +} + +// generationFromKubeContent extracts the hyperfleet.io/generation annotation +// from a previously-stored ApplyDesire's KubeContent. Returns 0 if the content +// can't be parsed or carries no annotation. +func generationFromKubeContent(kubeContent []byte) int64 { + obj, err := parseToUnstructured(kubeContent) + if err != nil { + return 0 + } + return manifest.GetGenerationFromUnstructured(obj) +} + +// parseToUnstructured parses JSON or YAML bytes into an unstructured resource, +// mirroring the pattern in internal/k8sclient/apply.go and +// internal/maestroclient/client.go. +func parseToUnstructured(data []byte) (*unstructured.Unstructured, error) { + obj := &unstructured.Unstructured{} + if err := json.Unmarshal(data, &obj.Object); err == nil && obj.Object != nil { + return obj, nil + } + jsonData, err := yaml.YAMLToJSON(data) + if err != nil { + return nil, fmt.Errorf("failed to convert YAML to JSON: %w", err) + } + if err := json.Unmarshal(jsonData, &obj.Object); err != nil { + return nil, fmt.Errorf("failed to parse manifest: %w", err) + } + return obj, nil +} diff --git a/internal/desireclient/types_test.go b/internal/desireclient/types_test.go new file mode 100644 index 00000000..c58732eb --- /dev/null +++ b/internal/desireclient/types_test.go @@ -0,0 +1,186 @@ +package desireclient + +import ( + "fmt" + "testing" + + "github.com/openshift-hyperfleet/hyperfleet-adapter/pkg/constants" + "github.com/openshift-hyperfleet/hyperfleet-applier/pkg/desire" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +// otherTransportContext stands in for a different transport's context type +// (e.g. *maestroclient.TransportContext) to exercise the type-assertion branch +// of resolveTransportContext without importing another transport package. +type otherTransportContext struct{} + +func TestResolveTransportContext_Invalid(t *testing.T) { + tests := []struct { + name string + target any + wantErr string + }{ + {name: "nil", target: nil, wantErr: "TransportContext"}, + {name: "wrong type", target: &otherTransportContext{}, wantErr: "TransportContext"}, + { + name: "missing management cluster", + target: &TransportContext{Resource: testResource}, + wantErr: "ManagementCluster", + }, + { + name: "missing resource", + target: &TransportContext{ManagementCluster: testManagementCluster}, + wantErr: "Resource", + }, + } + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + _, err := resolveTransportContext(tt.target) + require.Error(t, err) + assert.Contains(t, err.Error(), tt.wantErr) + }) + } +} + +func TestResolveTransportContext_Valid(t *testing.T) { + in := &TransportContext{ManagementCluster: testManagementCluster, Resource: testResource} + tc, err := resolveTransportContext(in) + require.NoError(t, err) + assert.Same(t, in, tc) +} + +func TestBuildIdentity(t *testing.T) { + tests := []struct { + name string + namespace string + resName string + wantErr string + }{ + {name: "valid", namespace: testNamespace, resName: testName}, + { + name: "invalid namespace is rejected", + namespace: "Invalid_Namespace", + resName: testName, + wantErr: "invalid identity", + }, + { + name: "invalid name is rejected", + namespace: testNamespace, + resName: "Invalid_Name", + wantErr: "invalid identity", + }, + } + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + id, err := buildIdentity(testTransportContext(), desire.TypeRead, testGVK(), tt.namespace, tt.resName) + if tt.wantErr != "" { + require.Error(t, err) + assert.Contains(t, err.Error(), tt.wantErr) + return + } + require.NoError(t, err) + assert.Equal(t, desire.Identity{ + ManagementCluster: testManagementCluster, + Type: desire.TypeRead, + Group: testGVK().Group, + Resource: testResource, + Namespace: tt.namespace, + Name: tt.resName, + }, id) + }) + } +} + +func TestParseToUnstructured(t *testing.T) { + tests := []struct { + name string + data []byte + wantErr bool + wantNil bool + }{ + {name: "valid JSON", data: configMapManifest(1)}, + { + name: "valid YAML", + data: []byte("apiVersion: v1\nkind: ConfigMap\nmetadata:\n name: my-config\n namespace: default\n"), + }, + {name: "garbage is neither valid JSON nor YAML", data: []byte("{not valid: [json or yaml"), wantErr: true}, + { + // json.Unmarshal("null", &obj.Object) succeeds with a nil map instead + // of erroring, and the YAML fallback degrades the same way — so a + // literal "null" manifest parses successfully into an empty object + // rather than failing loudly. Pinning this down since it's the one + // input where the two-stage JSON-then-YAML parse doesn't behave like + // either "valid" or "invalid". + name: "literal null parses to an empty object rather than erroring", + data: []byte("null"), + wantNil: true, + }, + } + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + obj, err := parseToUnstructured(tt.data) + if tt.wantErr { + require.Error(t, err) + return + } + require.NoError(t, err) + if tt.wantNil { + assert.Nil(t, obj.Object) + return + } + assert.Equal(t, "ConfigMap", obj.GetKind()) + assert.Equal(t, testName, obj.GetName()) + }) + } +} + +// TestParseToUnstructured_JSONAndYAMLProduceEquivalentObject verifies the +// property the store's JSON-only contract depends on: a JSON manifest and its +// YAML equivalent must parse to the identical unstructured object, so +// ApplyResource's re-marshal-to-JSON step doesn't silently diverge in content +// depending on which format the caller happened to render. +func TestParseToUnstructured_JSONAndYAMLProduceEquivalentObject(t *testing.T) { + jsonManifest := configMapManifest(1) + yamlManifest := fmt.Appendf(nil, ` +apiVersion: v1 +kind: ConfigMap +metadata: + name: %s + namespace: %s + annotations: + %s: "1" +data: + key: value +`, testName, testNamespace, constants.AnnotationGeneration) + + jsonObj, err := parseToUnstructured(jsonManifest) + require.NoError(t, err) + yamlObj, err := parseToUnstructured(yamlManifest) + require.NoError(t, err) + + assert.Equal(t, jsonObj.Object, yamlObj.Object, + "JSON and YAML encodings of the same manifest must parse to the same object") +} + +func TestGenerationFromKubeContent(t *testing.T) { + tests := []struct { + name string + content []byte + want int64 + }{ + {name: "valid content with generation annotation", content: configMapManifest(5), want: 5}, + { + name: "valid content without generation annotation", + content: []byte(`{"apiVersion":"v1","kind":"ConfigMap","metadata":{"name":"x","namespace":"default"}}`), + want: 0, + }, + {name: "unparseable content returns zero rather than erroring", content: []byte("not-json"), want: 0}, + {name: "empty content returns zero rather than erroring", content: nil, want: 0}, + } + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + assert.Equal(t, tt.want, generationFromKubeContent(tt.content)) + }) + } +}