From 8676a9d9afc1e628478e643954858054c63e8422 Mon Sep 17 00:00:00 2001 From: Michal Vavrinec Date: Thu, 17 Sep 2026 11:13:01 +0200 Subject: [PATCH 1/8] HYPERFLEET-1439 - feat: cleanup desires after confirmed resource deletion Add DesireCleaner interface and implement CleanupAfterDeletion on the desire client. When a resource is confirmed deleted (Step 2: already gone, Step 6: confirmed after delete), the executor cleans up the delete desire (only if Successful=True) and then the read desire. Scoped to by-name discovery only. Co-Authored-By: Claude Opus 4.6 --- internal/desireclient/cleanup.go | 75 +++++++ internal/desireclient/cleanup_test.go | 209 ++++++++++++++++++++ internal/desireclient/client.go | 5 +- internal/desireclient/helpers_test.go | 26 +++ internal/executor/resource_executor.go | 84 ++++++-- internal/executor/resource_executor_test.go | 168 ++++++++++++++++ internal/transportclient/interface.go | 14 ++ 7 files changed, 566 insertions(+), 15 deletions(-) create mode 100644 internal/desireclient/cleanup.go create mode 100644 internal/desireclient/cleanup_test.go diff --git a/internal/desireclient/cleanup.go b/internal/desireclient/cleanup.go new file mode 100644 index 00000000..c5bc472d --- /dev/null +++ b/internal/desireclient/cleanup.go @@ -0,0 +1,75 @@ +package desireclient + +import ( + "context" + "errors" + "fmt" + "log/slog" + + "github.com/openshift-hyperfleet/hyperfleet-adapter/internal/transportclient" + "github.com/openshift-hyperfleet/hyperfleet-applier/pkg/desire" + "k8s.io/apimachinery/pkg/runtime/schema" +) + +// CleanupAfterDeletion implements transportclient.DesireCleaner. It removes +// the delete desire (only when the applier confirms deletion) then the read +// desire. Returns an error if the delete desire exists but is not yet confirmed, +// causing the executor to retry on the next reconciliation. +func (c *Client) CleanupAfterDeletion( + ctx context.Context, + gvk schema.GroupVersionKind, + namespace, name string, + 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 + } + + dd, err := c.store.GetDeleteDesire(ctx, deleteID) + switch { + case errors.Is(err, desire.ErrNotFound): + // No delete desire — proceed to read desire cleanup. + case err != nil: + return fmt.Errorf("desireclient: cleanup: failed to get delete desire for %s/%s: %w", + namespace, name, err) + case !desire.IsDeleted(dd.Status): + return fmt.Errorf("desireclient: cleanup: deletion not yet confirmed for %s/%s", + namespace, name) + default: + if delErr := c.store.DeleteDeleteDesire(ctx, deleteID, c.owner, dd.Version); delErr != nil { + return fmt.Errorf("desireclient: cleanup: failed to delete delete desire for %s/%s: %w", + namespace, name, delErr) + } + slog.DebugContext(ctx, "desireclient: cleanup: removed confirmed delete desire", + "namespace", namespace, "name", name) + } + + readID, err := buildIdentity(tc, desire.TypeRead, gvk, namespace, name) + if err != nil { + return err + } + + rd, err := c.store.GetReadDesire(ctx, readID) + switch { + case errors.Is(err, desire.ErrNotFound): + return nil + case err != nil: + return fmt.Errorf("desireclient: cleanup: failed to get read desire for %s/%s: %w", + namespace, name, err) + default: + if delErr := c.store.DeleteReadDesire(ctx, readID, c.owner, rd.Version); delErr != nil { + return fmt.Errorf("desireclient: cleanup: failed to delete read desire for %s/%s: %w", + namespace, name, delErr) + } + slog.DebugContext(ctx, "desireclient: cleanup: removed read desire", + "namespace", namespace, "name", name) + } + + return nil +} diff --git a/internal/desireclient/cleanup_test.go b/internal/desireclient/cleanup_test.go new file mode 100644 index 00000000..4986acc0 --- /dev/null +++ b/internal/desireclient/cleanup_test.go @@ -0,0 +1,209 @@ +package desireclient + +import ( + "context" + "errors" + "testing" + + "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 TestCleanupAfterDeletion_ConfirmedDelete_RemovesBoth(t *testing.T) { + ctx := context.Background() + store := newMemoryStore() + c := newTestClient(store) + + deleteID := desire.Identity{ + ManagementCluster: testManagementCluster, Type: desire.TypeDelete, + Resource: testResource, Namespace: testNamespace, Name: testName, + } + readID := desire.Identity{ + ManagementCluster: testManagementCluster, Type: desire.TypeRead, + Resource: testResource, Namespace: testNamespace, Name: testName, + } + + putConfirmedDeleteDesire(t, ctx, store) + + _, err := store.CreateReadDesire(ctx, desire.ReadDesire{ + Identity: readID, Owner: testOwner, TargetVersion: "v1", + }) + require.NoError(t, err) + + err = c.CleanupAfterDeletion(ctx, testGVK(), testNamespace, testName, testTransportContext()) + require.NoError(t, err) + + _, err = store.GetDeleteDesire(ctx, deleteID) + assert.True(t, errors.Is(err, desire.ErrNotFound), "delete desire must be removed") + + _, err = store.GetReadDesire(ctx, readID) + assert.True(t, errors.Is(err, desire.ErrNotFound), "read desire must be removed") +} + +func TestCleanupAfterDeletion_PendingDelete_SkipsCleanup(t *testing.T) { + ctx := context.Background() + store := newMemoryStore() + c := newTestClient(store) + + deleteID := desire.Identity{ + ManagementCluster: testManagementCluster, Type: desire.TypeDelete, + Resource: testResource, Namespace: testNamespace, Name: testName, + } + readID := desire.Identity{ + ManagementCluster: testManagementCluster, Type: desire.TypeRead, + Resource: testResource, Namespace: testNamespace, Name: testName, + } + + putDeleteDesire(t, ctx, store, metav1.ConditionFalse, desire.ReasonWaitingForDeletion) + + _, err := store.CreateReadDesire(ctx, desire.ReadDesire{ + Identity: readID, Owner: testOwner, TargetVersion: "v1", + }) + require.NoError(t, err) + + err = c.CleanupAfterDeletion(ctx, testGVK(), testNamespace, testName, testTransportContext()) + require.Error(t, err, "pending delete desire must return an error") + assert.Contains(t, err.Error(), "deletion not yet confirmed") + + _, err = store.GetDeleteDesire(ctx, deleteID) + assert.NoError(t, err, "delete desire must still exist") + + _, err = store.GetReadDesire(ctx, readID) + assert.NoError(t, err, "read desire must still exist") +} + +func TestCleanupAfterDeletion_NoDeleteDesire_RemovesReadDesire(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, + } + _, err := store.CreateReadDesire(ctx, desire.ReadDesire{ + Identity: readID, Owner: testOwner, TargetVersion: "v1", + }) + require.NoError(t, err) + + err = c.CleanupAfterDeletion(ctx, testGVK(), testNamespace, testName, testTransportContext()) + require.NoError(t, err) + + _, err = store.GetReadDesire(ctx, readID) + assert.True(t, errors.Is(err, desire.ErrNotFound), "read desire must be removed") +} + +func TestCleanupAfterDeletion_NoDesires_NoError(t *testing.T) { + ctx := context.Background() + store := newMemoryStore() + c := newTestClient(store) + + err := c.CleanupAfterDeletion(ctx, testGVK(), testNamespace, testName, testTransportContext()) + require.NoError(t, err) +} + +func TestCleanupAfterDeletion_DeleteDesireOnly_NoReadDesire(t *testing.T) { + ctx := context.Background() + store := newMemoryStore() + c := newTestClient(store) + + deleteID := desire.Identity{ + ManagementCluster: testManagementCluster, Type: desire.TypeDelete, + Resource: testResource, Namespace: testNamespace, Name: testName, + } + + putConfirmedDeleteDesire(t, ctx, store) + + err := c.CleanupAfterDeletion(ctx, testGVK(), testNamespace, testName, testTransportContext()) + require.NoError(t, err) + + _, err = store.GetDeleteDesire(ctx, deleteID) + assert.True(t, errors.Is(err, desire.ErrNotFound), "delete desire must be removed") +} + +func TestCleanupAfterDeletion_RequiresTransportContext(t *testing.T) { + ctx := context.Background() + c := newTestClient(newMemoryStore()) + + err := c.CleanupAfterDeletion(ctx, testGVK(), testNamespace, testName, nil) + require.Error(t, err) +} + +func TestCleanupAfterDeletion_GetDeleteDesireError(t *testing.T) { + ctx := context.Background() + store := &failingGetDeleteDesireStore{SpecStore: newMemoryStore()} + c := newTestClient(store) + + err := c.CleanupAfterDeletion(ctx, testGVK(), testNamespace, testName, testTransportContext()) + require.Error(t, err) + assert.Contains(t, err.Error(), "failed to get delete desire") +} + +func TestCleanupAfterDeletion_DeleteDeleteDesireError(t *testing.T) { + ctx := context.Background() + inner := newMemoryStore() + + putConfirmedDeleteDesire(t, ctx, inner) + + store := &failingDeleteDeleteDesireStore{SpecStore: inner} + c := newTestClient(store) + + err := c.CleanupAfterDeletion(ctx, testGVK(), testNamespace, testName, testTransportContext()) + require.Error(t, err) + assert.Contains(t, err.Error(), "failed to delete delete desire") +} + +func TestCleanupAfterDeletion_DeleteReadDesireError(t *testing.T) { + ctx := context.Background() + inner := newMemoryStore() + + readID := desire.Identity{ + ManagementCluster: testManagementCluster, Type: desire.TypeRead, + Resource: testResource, Namespace: testNamespace, Name: testName, + } + _, err := inner.CreateReadDesire(ctx, desire.ReadDesire{ + Identity: readID, Owner: testOwner, TargetVersion: "v1", + }) + require.NoError(t, err) + + store := &failingDeleteReadDesireStore{SpecStore: inner} + c := newTestClient(store) + + err = c.CleanupAfterDeletion(ctx, testGVK(), testNamespace, testName, testTransportContext()) + require.Error(t, err) + assert.Contains(t, err.Error(), "failed to delete read desire") +} + +// --- Test store wrappers --- + +type failingGetDeleteDesireStore struct { + desire.SpecStore +} + +func (f *failingGetDeleteDesireStore) GetDeleteDesire( + _ context.Context, _ desire.Identity, +) (desire.DeleteDesire, error) { + return desire.DeleteDesire{}, errors.New("boom: store unavailable") +} + +type failingDeleteDeleteDesireStore struct { + desire.SpecStore +} + +func (f *failingDeleteDeleteDesireStore) DeleteDeleteDesire( + _ context.Context, _ desire.Identity, _ string, _ int64, +) error { + return errors.New("boom: version conflict") +} + +type failingDeleteReadDesireStore struct { + desire.SpecStore +} + +func (f *failingDeleteReadDesireStore) DeleteReadDesire( + _ context.Context, _ desire.Identity, _ string, _ int64, +) error { + return errors.New("boom: version conflict") +} diff --git a/internal/desireclient/client.go b/internal/desireclient/client.go index 6aa3a2cb..dbe579d7 100644 --- a/internal/desireclient/client.go +++ b/internal/desireclient/client.go @@ -26,4 +26,7 @@ func NewClient(store desire.SpecStore, owner string) *Client { return &Client{store: store, owner: owner} } -var _ transportclient.TransportClient = (*Client)(nil) +var ( + _ transportclient.TransportClient = (*Client)(nil) + _ transportclient.DesireCleaner = (*Client)(nil) +) diff --git a/internal/desireclient/helpers_test.go b/internal/desireclient/helpers_test.go index e5768b13..3ac14340 100644 --- a/internal/desireclient/helpers_test.go +++ b/internal/desireclient/helpers_test.go @@ -82,6 +82,32 @@ func putKubeAPIErrorReadDesire( }) } +// putDeleteDesire creates a delete desire with the given condition status and reason. +func putDeleteDesire( + t *testing.T, ctx context.Context, store *memory.Store, + condStatus metav1.ConditionStatus, reason string, +) { + t.Helper() + id := desire.Identity{ + ManagementCluster: testManagementCluster, Type: desire.TypeDelete, + Resource: testResource, Namespace: testNamespace, Name: testName, + } + dd, err := store.CreateDeleteDesire(ctx, desire.DeleteDesire{Identity: id, Owner: testOwner}) + require.NoError(t, err) + _, err = store.UpdateDeleteDesireStatus(ctx, id, desire.Status{ + Conditions: []metav1.Condition{{ + Type: desire.TypeSuccessful, Status: condStatus, Reason: reason, + }}, + }, dd.Version) + require.NoError(t, err) +} + +// putConfirmedDeleteDesire creates a delete desire marked as successfully deleted. +func putConfirmedDeleteDesire(t *testing.T, ctx context.Context, store *memory.Store) { + t.Helper() + putDeleteDesire(t, ctx, store, metav1.ConditionTrue, desire.ReasonDeleted) +} + // 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} diff --git a/internal/executor/resource_executor.go b/internal/executor/resource_executor.go index 5186849a..144de7e0 100644 --- a/internal/executor/resource_executor.go +++ b/internal/executor/resource_executor.go @@ -291,6 +291,30 @@ func (re *ResourceExecutor) renderToBytes( // For k8s transport: discovers the K8s resource by name or label selector. // For maestro transport: discovers the ManifestWork by name or label selector. // The discovered resource is stored in execCtx.Resources for post-action CEL evaluation. +type discoveryTarget struct { + Namespace string + Name string +} + +func (re *ResourceExecutor) renderDiscoveryTarget( + discovery *configloader.DiscoveryConfig, + params map[string]any, +) (*discoveryTarget, error) { + namespace, err := utils.RenderTemplate(discovery.Namespace, params) + if err != nil { + return nil, fmt.Errorf("failed to render namespace template: %w", err) + } + dt := &discoveryTarget{Namespace: namespace} + if discovery.ByName != "" { + name, err := utils.RenderTemplate(discovery.ByName, params) + if err != nil { + return nil, fmt.Errorf("failed to render byName template: %w", err) + } + dt.Name = name + } + return dt, nil +} + func (re *ResourceExecutor) discoverResource( ctx context.Context, resource configloader.Resource, @@ -303,24 +327,15 @@ func (re *ResourceExecutor) discoverResource( return nil, nil } - // Render discovery namespace template - namespace, err := utils.RenderTemplate(discovery.Namespace, execCtx.Params) + dt, err := re.renderDiscoveryTarget(discovery, execCtx.Params) if err != nil { - return nil, fmt.Errorf("failed to render namespace template: %w", err) + return nil, err } // Discover by name if discovery.ByName != "" { - name, err := utils.RenderTemplate(discovery.ByName, execCtx.Params) - if err != nil { - return nil, fmt.Errorf("failed to render byName template: %w", err) - } - - // For maestro: use ManifestWork GVK - // For k8s: parse the rendered manifest to get GVK gvk := re.resolveGVK(resource) - - return transportClient.GetResource(ctx, gvk, namespace, name, transportTarget) + return transportClient.GetResource(ctx, gvk, dt.Namespace, dt.Name, transportTarget) } // Discover by label selector @@ -340,7 +355,7 @@ func (re *ResourceExecutor) discoverResource( labelSelector := manifest.BuildLabelSelector(renderedLabels) discoveryConfig := &manifest.DiscoveryConfig{ - Namespace: namespace, + Namespace: dt.Namespace, LabelSelector: labelSelector, } @@ -361,6 +376,25 @@ func (re *ResourceExecutor) discoverResource( return nil, fmt.Errorf("discovery config must specify byName or bySelectors") } +func (re *ResourceExecutor) tryCleanupDesires( + ctx context.Context, + resource configloader.Resource, + execCtx *ExecutionContext, + transportClient transportclient.TransportClient, + transportTarget transportclient.TransportContext, + gvk schema.GroupVersionKind, +) error { + cleaner, ok := transportClient.(transportclient.DesireCleaner) + if !ok || resource.Discovery == nil || resource.Discovery.ByName == "" { + return nil + } + dt, err := re.renderDiscoveryTarget(resource.Discovery, execCtx.Params) + if err != nil { + return err + } + return cleaner.CleanupAfterDeletion(ctx, gvk, dt.Namespace, dt.Name, transportTarget) +} + // discoverNestedResources discovers sub-resources within a parent resource (e.g., manifests inside a ManifestWork). // Each nestedDiscovery is matched against the parent's nested manifests using manifest.DiscoverNestedManifest. func (re *ResourceExecutor) discoverNestedResources( @@ -685,8 +719,20 @@ func (re *ResourceExecutor) executeResourceDelete( // Store nil — the key is removed from the CEL resources map, so // !resources.?X.hasValue() evaluates to true in this reconciliation. execCtx.Resources[resource.Name] = nil - result.OperationReason = "resource already deleted or never existed" + + if err := re.tryCleanupDesires(ctx, resource, execCtx, transportClient, transportTarget, gvk); err != nil { + slog.ErrorContext(ctx, "resource desire cleanup failed after delete", + "resource", resource.Name, "error", err) + result.Status = StatusFailed + result.Error = err + re.recordResourceError(execCtx, resource, err) + re.metrics.RecordDeletion(resourceType, metrics.DeletionStatusError) + re.metrics.ObserveDeletionDuration(resourceType, time.Since(startTime)) + return result, NewExecutorError(PhaseResources, resource.Name, "desire cleanup failed", err) + } + slog.InfoContext(ctx, "resource delete: already deleted or never existed", "resource", resource.Name) + result.OperationReason = "resource already deleted or never existed" re.metrics.RecordDeletion(resourceType, metrics.DeletionStatusSuccess) re.metrics.ObserveDeletionDuration(resourceType, time.Since(startTime)) return result, nil @@ -744,6 +790,16 @@ func (re *ResourceExecutor) executeResourceDelete( // Resource is confirmed gone: dependent resources can proceed in this reconciliation. execCtx.Resources[resource.Name] = nil slog.DebugContext(ctx, "resource confirmed deleted (post-delete discovery: not found)", "resource", resource.Name) + if err := re.tryCleanupDesires(ctx, resource, execCtx, transportClient, transportTarget, gvk); err != nil { + slog.ErrorContext(ctx, "resource desire cleanup failed after delete", + "resource", resource.Name, "error", err) + result.Status = StatusFailed + result.Error = err + re.recordResourceError(execCtx, resource, err) + re.metrics.RecordDeletion(resourceType, metrics.DeletionStatusError) + re.metrics.ObserveDeletionDuration(resourceType, time.Since(startTime)) + return result, NewExecutorError(PhaseResources, resource.Name, "desire cleanup failed", err) + } default: // Resource still present (finalizers or async deletion): dependents wait for next reconciliation. execCtx.Resources[resource.Name] = postDeleteDiscovered diff --git a/internal/executor/resource_executor_test.go b/internal/executor/resource_executor_test.go index 11d19bac..ee593861 100644 --- a/internal/executor/resource_executor_test.go +++ b/internal/executor/resource_executor_test.go @@ -2284,3 +2284,171 @@ func TestResourceExecutor_LifecycleDelete_BySelectors(t *testing.T) { assert.True(t, exists, "nil sentinel should be in execCtx.Resources") assert.Nil(t, storedVal, "nil stored when post-delete discovery finds no resources") } + +// ---- DesireCleaner integration ---- + +func TestResourceExecutor_LifecycleDelete_Step2_CleanupCalled(t *testing.T) { + inner := k8sclient.NewMockK8sClient() + mock := &cleanupTrackingDeleteMockClient{ + MockK8sClient: inner, + } + + re := newResourceExecutor(&ExecutorConfig{ + TransportRegistry: testTransportRegistry(mock), + }) + + resource := newResourceWithLifecycle("deleted_time != null", "Background") + execCtx := NewExecutionContext(context.Background(), nil, nil) + execCtx.Params["deleted_time"] = testDeletedTime + + results, err := re.ExecuteAll(context.Background(), []configloader.Resource{resource}, execCtx) + + require.NoError(t, err) + require.Len(t, results, 1) + assert.Equal(t, StatusSuccess, results[0].Status) + assert.False(t, mock.DeleteCalled, "DeleteResource must not be called when resource was already gone") + assert.True(t, mock.CleanupCalled, "CleanupAfterDeletion must be called when resource is not found") + assert.Equal(t, "default", mock.CleanupNamespace) + assert.Equal(t, "test-cm", mock.CleanupName) +} + +func TestResourceExecutor_LifecycleDelete_Step6_CleanupCalled(t *testing.T) { + discovered := &unstructured.Unstructured{ + Object: map[string]interface{}{ + "apiVersion": "v1", + "kind": "ConfigMap", + "metadata": map[string]interface{}{"name": "test-cm", "namespace": "default"}, + }, + } + + inner := k8sclient.NewMockK8sClient() + // Resource exists initially, DeleteResource removes it from Resources map, + // so post-delete GetResource returns NotFound. + inner.Resources["default/test-cm"] = discovered + mock := &cleanupTrackingDeleteMockClient{ + MockK8sClient: inner, + } + + re := newResourceExecutor(&ExecutorConfig{ + TransportRegistry: testTransportRegistry(mock), + }) + + resource := newResourceWithLifecycle("deleted_time != null", "Background") + execCtx := NewExecutionContext(context.Background(), nil, nil) + execCtx.Params["deleted_time"] = testDeletedTime + + results, err := re.ExecuteAll(context.Background(), []configloader.Resource{resource}, execCtx) + + require.NoError(t, err) + require.Len(t, results, 1) + assert.Equal(t, StatusSuccess, results[0].Status) + assert.True(t, mock.DeleteCalled, "DeleteResource must be called") + assert.True(t, mock.CleanupCalled, "CleanupAfterDeletion must be called when post-delete discovery confirms gone") + assert.Equal(t, "default", mock.CleanupNamespace) + assert.Equal(t, "test-cm", mock.CleanupName) +} + +// cleanupTrackingDeleteMockClient delegates to MockK8sClient (which removes resources on delete) +// and implements DesireCleaner to track cleanup calls. +type cleanupTrackingDeleteMockClient struct { + *k8sclient.MockK8sClient + CleanupError error + CleanupNamespace string + CleanupName string + DeleteCalled bool + CleanupCalled bool +} + +func (m *cleanupTrackingDeleteMockClient) DeleteResource( + ctx context.Context, + gvk schema.GroupVersionKind, + namespace, name string, + opts *transportclient.DeleteOptions, + target transportclient.TransportContext, +) error { + m.DeleteCalled = true + return m.MockK8sClient.DeleteResource(ctx, gvk, namespace, name, opts, target) +} + +func (m *cleanupTrackingDeleteMockClient) CleanupAfterDeletion( + _ context.Context, + _ schema.GroupVersionKind, + namespace, name string, + _ transportclient.TransportContext, +) error { + m.CleanupCalled = true + m.CleanupNamespace = namespace + m.CleanupName = name + return m.CleanupError +} + +func TestResourceExecutor_LifecycleDelete_CleanupError_StatusFailed(t *testing.T) { + inner := k8sclient.NewMockK8sClient() + mock := &cleanupTrackingDeleteMockClient{ + MockK8sClient: inner, + CleanupError: errors.New("cleanup: deletion not yet confirmed"), + } + + re := newResourceExecutor(&ExecutorConfig{ + TransportRegistry: testTransportRegistry(mock), + }) + + resource := newResourceWithLifecycle("deleted_time != null", "Background") + execCtx := NewExecutionContext(context.Background(), nil, nil) + execCtx.Params["deleted_time"] = testDeletedTime + + results, err := re.ExecuteAll(context.Background(), []configloader.Resource{resource}, execCtx) + + require.Error(t, err) + require.Len(t, results, 1) + assert.Equal(t, StatusFailed, results[0].Status) + assert.True(t, mock.CleanupCalled) +} + +// cleanupKeepOnDeleteMockClient keeps the resource after delete (simulating finalizers/async) +// and implements DesireCleaner to track whether cleanup was attempted. +type cleanupKeepOnDeleteMockClient struct { + *keepOnDeleteMockClient + CleanupCalled bool +} + +func (m *cleanupKeepOnDeleteMockClient) CleanupAfterDeletion( + _ context.Context, + _ schema.GroupVersionKind, + _, _ string, + _ transportclient.TransportContext, +) error { + m.CleanupCalled = true + return nil +} + +func TestResourceExecutor_LifecycleDelete_StillPresent_NoCleanup(t *testing.T) { + discovered := &unstructured.Unstructured{ + Object: map[string]interface{}{ + "apiVersion": "v1", + "kind": "ConfigMap", + "metadata": map[string]interface{}{"name": "test-cm", "namespace": "default"}, + }, + } + + inner := k8sclient.NewMockK8sClient() + inner.Resources["default/test-cm"] = discovered + mock := &cleanupKeepOnDeleteMockClient{ + keepOnDeleteMockClient: &keepOnDeleteMockClient{MockK8sClient: inner}, + } + + re := newResourceExecutor(&ExecutorConfig{ + TransportRegistry: testTransportRegistry(mock), + }) + + resource := newResourceWithLifecycle("deleted_time != null", "Background") + execCtx := NewExecutionContext(context.Background(), nil, nil) + execCtx.Params["deleted_time"] = testDeletedTime + + results, err := re.ExecuteAll(context.Background(), []configloader.Resource{resource}, execCtx) + + require.NoError(t, err) + require.Len(t, results, 1) + assert.Equal(t, StatusSuccess, results[0].Status) + assert.False(t, mock.CleanupCalled, "CleanupAfterDeletion must not be called when resource is still present") +} diff --git a/internal/transportclient/interface.go b/internal/transportclient/interface.go index 20b8e306..cc49bc05 100644 --- a/internal/transportclient/interface.go +++ b/internal/transportclient/interface.go @@ -71,3 +71,17 @@ type TransportClient interface { target TransportContext, ) error } + +// DesireCleaner is an optional interface that transport clients may +// implement to remove transport-layer bookkeeping after a resource has been +// confirmed deleted. The executor calls this when pre-delete discovery +// confirms the resource is already gone, and when post-delete re-discovery +// confirms it was removed. +type DesireCleaner interface { + CleanupAfterDeletion( + ctx context.Context, + gvk schema.GroupVersionKind, + namespace, name string, + target TransportContext, + ) error +} From a4737999a4713b91dab5bfd29e82993253a43495 Mon Sep 17 00:00:00 2001 From: Michal Vavrinec Date: Fri, 18 Sep 2026 12:13:53 +0200 Subject: [PATCH 2/8] HYPERFLEET-1439 - refactor: migrate test helpers to exported testhelpers.go Replace unexported helpers_test.go with exported TestIdentity builder and shared helpers in testhelpers.go, enabling reuse from other packages. Co-Authored-By: Claude Opus 4.6 --- internal/desireclient/cleanup_test.go | 45 ++++- internal/desireclient/desireclient_test.go | 48 ++++++ internal/desireclient/discover_test.go | 26 +-- internal/desireclient/get_test.go | 31 ++-- internal/desireclient/helpers_test.go | 155 ----------------- internal/desireclient/testhelpers.go | 185 +++++++++++++++++++++ 6 files changed, 299 insertions(+), 191 deletions(-) delete mode 100644 internal/desireclient/helpers_test.go create mode 100644 internal/desireclient/testhelpers.go diff --git a/internal/desireclient/cleanup_test.go b/internal/desireclient/cleanup_test.go index 4986acc0..0ff9bce3 100644 --- a/internal/desireclient/cleanup_test.go +++ b/internal/desireclient/cleanup_test.go @@ -25,7 +25,7 @@ func TestCleanupAfterDeletion_ConfirmedDelete_RemovesBoth(t *testing.T) { Resource: testResource, Namespace: testNamespace, Name: testName, } - putConfirmedDeleteDesire(t, ctx, store) + PutConfirmedDeleteDesire(t, ctx, store, testID.Delete(), testOwner) _, err := store.CreateReadDesire(ctx, desire.ReadDesire{ Identity: readID, Owner: testOwner, TargetVersion: "v1", @@ -56,7 +56,7 @@ func TestCleanupAfterDeletion_PendingDelete_SkipsCleanup(t *testing.T) { Resource: testResource, Namespace: testNamespace, Name: testName, } - putDeleteDesire(t, ctx, store, metav1.ConditionFalse, desire.ReasonWaitingForDeletion) + PutDeleteDesire(t, ctx, store, testID.Delete(), testOwner, metav1.ConditionFalse, desire.ReasonWaitingForDeletion) _, err := store.CreateReadDesire(ctx, desire.ReadDesire{ Identity: readID, Owner: testOwner, TargetVersion: "v1", @@ -66,6 +66,7 @@ func TestCleanupAfterDeletion_PendingDelete_SkipsCleanup(t *testing.T) { err = c.CleanupAfterDeletion(ctx, testGVK(), testNamespace, testName, testTransportContext()) require.Error(t, err, "pending delete desire must return an error") assert.Contains(t, err.Error(), "deletion not yet confirmed") + assert.True(t, errors.Is(err, ErrDeletionPending), "must wrap ErrDeletionPending") _, err = store.GetDeleteDesire(ctx, deleteID) assert.NoError(t, err, "delete desire must still exist") @@ -104,6 +105,42 @@ func TestCleanupAfterDeletion_NoDesires_NoError(t *testing.T) { require.NoError(t, err) } +func TestCleanupAfterDeletion_ApplyDesireExists_NoDeleteDesire_ReturnsError(t *testing.T) { + ctx := context.Background() + store := newMemoryStore() + c := newTestClient(store) + + applyID := desire.Identity{ + ManagementCluster: testManagementCluster, Type: desire.TypeApply, + Resource: testResource, Namespace: testNamespace, Name: testName, + } + _, err := store.CreateApplyDesire(ctx, desire.ApplyDesire{ + Identity: applyID, Owner: testOwner, + Spec: desire.ApplySpec{KubeContent: configMapManifest(1)}, + }) + require.NoError(t, err) + + readID := desire.Identity{ + ManagementCluster: testManagementCluster, Type: desire.TypeRead, + Resource: testResource, Namespace: testNamespace, Name: testName, + } + _, err = store.CreateReadDesire(ctx, desire.ReadDesire{ + Identity: readID, Owner: testOwner, TargetVersion: "v1", + }) + require.NoError(t, err) + + err = c.CleanupAfterDeletion(ctx, testGVK(), testNamespace, testName, testTransportContext()) + require.Error(t, err) + assert.Contains(t, err.Error(), "apply desire still exists") + assert.True(t, errors.Is(err, ErrDeletionPending), "must wrap ErrDeletionPending") + + _, err = store.GetApplyDesire(ctx, applyID) + assert.NoError(t, err, "apply desire must still exist") + + _, err = store.GetReadDesire(ctx, readID) + assert.NoError(t, err, "read desire must still exist") +} + func TestCleanupAfterDeletion_DeleteDesireOnly_NoReadDesire(t *testing.T) { ctx := context.Background() store := newMemoryStore() @@ -114,7 +151,7 @@ func TestCleanupAfterDeletion_DeleteDesireOnly_NoReadDesire(t *testing.T) { Resource: testResource, Namespace: testNamespace, Name: testName, } - putConfirmedDeleteDesire(t, ctx, store) + PutConfirmedDeleteDesire(t, ctx, store, testID.Delete(), testOwner) err := c.CleanupAfterDeletion(ctx, testGVK(), testNamespace, testName, testTransportContext()) require.NoError(t, err) @@ -145,7 +182,7 @@ func TestCleanupAfterDeletion_DeleteDeleteDesireError(t *testing.T) { ctx := context.Background() inner := newMemoryStore() - putConfirmedDeleteDesire(t, ctx, inner) + PutConfirmedDeleteDesire(t, ctx, inner, testID.Delete(), testOwner) store := &failingDeleteDeleteDesireStore{SpecStore: inner} c := newTestClient(store) diff --git a/internal/desireclient/desireclient_test.go b/internal/desireclient/desireclient_test.go index 80f8e774..b837037d 100644 --- a/internal/desireclient/desireclient_test.go +++ b/internal/desireclient/desireclient_test.go @@ -72,6 +72,54 @@ func (f *failingCreateDeleteDesireStore) CreateDeleteDesire( return desire.DeleteDesire{}, errors.New("boom: delete desire store unavailable") } +var testID = TestIdentity{ + ManagementCluster: testManagementCluster, + Resource: testResource, + Namespace: testNamespace, + Name: testName, +} + func newMemoryStore() *memory.Store { return memory.New() } + +// 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/discover_test.go b/internal/desireclient/discover_test.go index f5911c68..add9d75b 100644 --- a/internal/desireclient/discover_test.go +++ b/internal/desireclient/discover_test.go @@ -25,7 +25,7 @@ func TestDiscoverResources_ReturnsSyncedResourceByName(t *testing.T) { store := newMemoryStore() c := newTestClient(store) - putSyncedReadDesire(t, ctx, store, testNamespace, testName, configMapManifest(1)) + PutSyncedReadDesire(t, ctx, store, testID.Read(), testOwner, configMapManifest(1)) list, err := c.DiscoverResources(ctx, testGVK(), &manifest.DiscoveryConfig{ByName: testName}, testTransportContext()) require.NoError(t, err) @@ -38,7 +38,7 @@ func TestDiscoverResources_ByNameExcludesNonMatchingName(t *testing.T) { store := newMemoryStore() c := newTestClient(store) - putSyncedReadDesire(t, ctx, store, testNamespace, testName, configMapManifest(1)) + PutSyncedReadDesire(t, ctx, store, testID.Read(), testOwner, configMapManifest(1)) discovery := &manifest.DiscoveryConfig{ByName: "other-name"} list, err := c.DiscoverResources(ctx, testGVK(), discovery, testTransportContext()) @@ -59,8 +59,8 @@ func TestDiscoverResources_LabelSelectorMatchesSubset(t *testing.T) { "apiVersion": "v1", "kind": "ConfigMap", "metadata": {"name": "unlabeled", "namespace": "default"} }`) - putSyncedReadDesire(t, ctx, store, "labeled", "labeled", labeledManifest) - putSyncedReadDesire(t, ctx, store, "unlabeled", "unlabeled", unlabeledManifest) + PutSyncedReadDesire(t, ctx, store, testID.WithNamespace("labeled").WithName("labeled").Read(), testOwner, labeledManifest) + PutSyncedReadDesire(t, ctx, store, testID.WithNamespace("unlabeled").WithName("unlabeled").Read(), testOwner, unlabeledManifest) list, err := c.DiscoverResources(ctx, testGVK(), &manifest.DiscoveryConfig{LabelSelector: "app=myapp"}, testTransportContext()) @@ -91,7 +91,7 @@ func TestDiscoverResources_SurfacesRetainedMirrorOnFailedRead(t *testing.T) { store := newMemoryStore() c := newTestClient(store) - putKubeAPIErrorReadDesire(t, ctx, store, testNamespace, testName, configMapManifest(1)) + PutKubeAPIErrorReadDesire(t, ctx, store, testID.Read(), testOwner, configMapManifest(1)) list, err := c.DiscoverResources(ctx, testGVK(), &manifest.DiscoveryConfig{}, testTransportContext()) require.NoError(t, err) @@ -105,7 +105,7 @@ func TestDiscoverResources_SkipsFailedReadWithNoRetainedMirror(t *testing.T) { store := newMemoryStore() c := newTestClient(store) - putKubeAPIErrorReadDesire(t, ctx, store, testNamespace, testName, nil) + PutKubeAPIErrorReadDesire(t, ctx, store, testID.Read(), testOwner, nil) list, err := c.DiscoverResources(ctx, testGVK(), &manifest.DiscoveryConfig{}, testTransportContext()) require.NoError(t, err) @@ -117,7 +117,7 @@ func TestDiscoverResources_SkipsNotFoundFalseDesire(t *testing.T) { store := newMemoryStore() c := newTestClient(store) - putNotFoundReadDesire(t, ctx, store, testNamespace, testName) + PutNotFoundReadDesire(t, ctx, store, testID.Read(), testOwner) list, err := c.DiscoverResources(ctx, testGVK(), &manifest.DiscoveryConfig{}, testTransportContext()) require.NoError(t, err) @@ -129,8 +129,8 @@ func TestDiscoverResources_SkipsUndecodableContentButKeepsOthers(t *testing.T) { store := newMemoryStore() c := newTestClient(store) - putSyncedReadDesire(t, ctx, store, testNamespace, "bad", []byte("not-json")) - putSyncedReadDesire(t, ctx, store, testNamespace, "good", configMapManifest(1)) + PutSyncedReadDesire(t, ctx, store, testID.WithName("bad").Read(), testOwner, []byte("not-json")) + PutSyncedReadDesire(t, ctx, store, testID.WithName("good").Read(), testOwner, 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") @@ -143,7 +143,7 @@ func TestDiscoverResources_FiltersOutOtherResourceType(t *testing.T) { store := newMemoryStore() c := newTestClient(store) - putSyncedReadDesire(t, ctx, store, testNamespace, testName, configMapManifest(1)) + PutSyncedReadDesire(t, ctx, store, testID.Read(), testOwner, configMapManifest(1)) otherContext := &TransportContext{ManagementCluster: testManagementCluster, Resource: "secrets"} list, err := c.DiscoverResources(ctx, testGVK(), &manifest.DiscoveryConfig{}, otherContext) @@ -182,7 +182,7 @@ func TestDiscoverResources_ScopedToPartition(t *testing.T) { store := newMemoryStore() c := newTestClient(store) - putSyncedReadDesire(t, ctx, store, testNamespace, testName, configMapManifest(1)) + PutSyncedReadDesire(t, ctx, store, testID.Read(), testOwner, configMapManifest(1)) otherPartition := &TransportContext{ManagementCluster: "other-cluster", Resource: testResource} list, err := c.DiscoverResources(ctx, testGVK(), &manifest.DiscoveryConfig{}, otherPartition) @@ -197,8 +197,8 @@ func TestDiscoverResources_MultipleMatches(t *testing.T) { 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) + PutSyncedReadDesire(t, ctx, store, testID.WithName("first").Read(), testOwner, first) + PutSyncedReadDesire(t, ctx, store, testID.WithName("second").Read(), testOwner, second) list, err := c.DiscoverResources(ctx, testGVK(), &manifest.DiscoveryConfig{}, testTransportContext()) require.NoError(t, err) diff --git a/internal/desireclient/get_test.go b/internal/desireclient/get_test.go index 1804aa81..bbc3fa5f 100644 --- a/internal/desireclient/get_test.go +++ b/internal/desireclient/get_test.go @@ -12,13 +12,6 @@ import ( 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()) @@ -35,7 +28,7 @@ func TestGetResource_ReadDesireExistsNoResourceObservedIsNotSyncedYet(t *testing c := newTestClient(store) _, err := store.CreateReadDesire(ctx, desire.ReadDesire{ - Identity: readIdentity(), Owner: testOwner, TargetVersion: "v1", + Identity: testID.Read(), Owner: testOwner, TargetVersion: "v1", }) require.NoError(t, err) @@ -49,7 +42,7 @@ func TestGetResource_SyncedReturnsMirroredObject(t *testing.T) { store := newMemoryStore() c := newTestClient(store) - putSyncedReadDesire(t, ctx, store, testNamespace, testName, configMapManifest(1)) + PutSyncedReadDesire(t, ctx, store, testID.Read(), testOwner, configMapManifest(1)) obj, err := c.GetResource(ctx, testGVK(), testNamespace, testName, testTransportContext()) require.NoError(t, err) @@ -62,7 +55,7 @@ func TestGetResource_ConfirmedNotFound(t *testing.T) { store := newMemoryStore() c := newTestClient(store) - putConfirmedAbsentReadDesire(t, ctx, store, testNamespace, testName) + PutConfirmedAbsentReadDesire(t, ctx, store, testID.Read(), testOwner) _, err := c.GetResource(ctx, testGVK(), testNamespace, testName, testTransportContext()) require.Error(t, err) @@ -75,7 +68,7 @@ func TestGetResource_InvalidReadDesire(t *testing.T) { store := newMemoryStore() c := newTestClient(store) - putInvalidReadDesire(t, ctx, store, testNamespace, testName) + PutInvalidReadDesire(t, ctx, store, testID.Read(), testOwner) _, err := c.GetResource(ctx, testGVK(), testNamespace, testName, testTransportContext()) require.Error(t, err) @@ -90,7 +83,7 @@ func TestGetResource_K8sAPIErrorWithRetainedMirrorReturnsStaleContent(t *testing store := newMemoryStore() c := newTestClient(store) - putKubeAPIErrorReadDesire(t, ctx, store, testNamespace, testName, configMapManifest(1)) + PutKubeAPIErrorReadDesire(t, ctx, store, testID.Read(), testOwner, configMapManifest(1)) obj, err := c.GetResource(ctx, testGVK(), testNamespace, testName, testTransportContext()) require.NoError(t, err, @@ -104,7 +97,7 @@ func TestGetResource_K8sAPIErrorWithNoMirrorYetIsNotSyncedYet(t *testing.T) { store := newMemoryStore() c := newTestClient(store) - putKubeAPIErrorReadDesire(t, ctx, store, testNamespace, testName, nil) + PutKubeAPIErrorReadDesire(t, ctx, store, testID.Read(), testOwner, nil) _, err := c.GetResource(ctx, testGVK(), testNamespace, testName, testTransportContext()) require.Error(t, err) @@ -159,30 +152,30 @@ func TestDecodeReadDesire(t *testing.T) { }, { name: "successful true decodes content", - condition: successfulCondition(metav1.ConditionTrue, desire.ReasonSynced), + condition: SuccessfulCondition(metav1.ConditionTrue, desire.ReasonSynced), content: configMapManifest(1), wantContent: true, }, { // ConditionTrue/ReasonNotFound is a shape the applier never reports; with no content it reads as not-synced-yet. name: "successful true with empty content is not synced yet", - condition: successfulCondition(metav1.ConditionTrue, desire.ReasonNotFound), + condition: SuccessfulCondition(metav1.ConditionTrue, desire.ReasonNotFound), wantNotSynced: true, }, { name: "false with notfound reason is confirmed absent", - condition: successfulCondition(metav1.ConditionFalse, desire.ReasonNotFound), + 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), + condition: SuccessfulCondition(metav1.ConditionFalse, desire.ReasonKubeAPIError), content: configMapManifest(1), wantContent: true, }, { name: "false with other reason and no retained mirror is not synced yet", - condition: successfulCondition(metav1.ConditionFalse, desire.ReasonKubeAPIError), + condition: SuccessfulCondition(metav1.ConditionFalse, desire.ReasonKubeAPIError), wantNotSynced: true, }, } @@ -190,7 +183,7 @@ func TestDecodeReadDesire(t *testing.T) { c := newTestClient(newMemoryStore()) for _, tt := range tests { t.Run(tt.name, func(t *testing.T) { - rd := desire.ReadDesire{Identity: readIdentity()} + rd := desire.ReadDesire{Identity: testID.Read()} if tt.condition != nil { rd.Status = desire.ReadStatus{ Status: desire.Status{Conditions: []metav1.Condition{*tt.condition}}, diff --git a/internal/desireclient/helpers_test.go b/internal/desireclient/helpers_test.go deleted file mode 100644 index 3ac14340..00000000 --- a/internal/desireclient/helpers_test.go +++ /dev/null @@ -1,155 +0,0 @@ -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, - }) -} - -// putDeleteDesire creates a delete desire with the given condition status and reason. -func putDeleteDesire( - t *testing.T, ctx context.Context, store *memory.Store, - condStatus metav1.ConditionStatus, reason string, -) { - t.Helper() - id := desire.Identity{ - ManagementCluster: testManagementCluster, Type: desire.TypeDelete, - Resource: testResource, Namespace: testNamespace, Name: testName, - } - dd, err := store.CreateDeleteDesire(ctx, desire.DeleteDesire{Identity: id, Owner: testOwner}) - require.NoError(t, err) - _, err = store.UpdateDeleteDesireStatus(ctx, id, desire.Status{ - Conditions: []metav1.Condition{{ - Type: desire.TypeSuccessful, Status: condStatus, Reason: reason, - }}, - }, dd.Version) - require.NoError(t, err) -} - -// putConfirmedDeleteDesire creates a delete desire marked as successfully deleted. -func putConfirmedDeleteDesire(t *testing.T, ctx context.Context, store *memory.Store) { - t.Helper() - putDeleteDesire(t, ctx, store, metav1.ConditionTrue, desire.ReasonDeleted) -} - -// 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/testhelpers.go b/internal/desireclient/testhelpers.go new file mode 100644 index 00000000..044c6a0d --- /dev/null +++ b/internal/desireclient/testhelpers.go @@ -0,0 +1,185 @@ +package desireclient + +import ( + "context" + "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 identity, owner, and status. +func PutReadDesire( + t testing.TB, ctx context.Context, store *memory.Store, + id desire.Identity, owner string, status desire.ReadStatus, +) { + t.Helper() + _, err := store.CreateReadDesire(ctx, desire.ReadDesire{ + Identity: id, Owner: owner, TargetVersion: "v1", + }) + require.NoError(t, err) + _, err = store.UpdateReadDesireStatus(ctx, id, status) + require.NoError(t, err) +} + +// PutNotFoundReadDesire creates a read desire marked as not found. +func PutNotFoundReadDesire( + t testing.TB, ctx context.Context, store *memory.Store, + id desire.Identity, owner string, +) { + t.Helper() + PutReadDesire(t, ctx, store, id, owner, desire.ReadStatus{ + Status: desire.Status{Conditions: []metav1.Condition{{ + Type: desire.TypeSuccessful, Status: metav1.ConditionFalse, Reason: desire.ReasonNotFound, + }}}, + }) +} + +// PutConfirmedAbsentReadDesire creates a read desire marked as not found. +func PutConfirmedAbsentReadDesire( + t testing.TB, ctx context.Context, store *memory.Store, + id desire.Identity, owner string, +) { + t.Helper() + PutNotFoundReadDesire(t, ctx, store, id, owner) +} + +// PutSyncedReadDesire creates a read desire with synced content. +func PutSyncedReadDesire( + t testing.TB, ctx context.Context, store *memory.Store, + id desire.Identity, owner string, content []byte, +) { + t.Helper() + PutReadDesire(t, ctx, store, id, owner, desire.ReadStatus{ + Status: desire.Status{Conditions: []metav1.Condition{{ + Type: desire.TypeSuccessful, Status: metav1.ConditionTrue, Reason: desire.ReasonSynced, + }}}, + KubeContent: content, + }) +} + +// PutInvalidReadDesire creates a read desire with Successful=True but Reason=NotFound. +func PutInvalidReadDesire( + t testing.TB, ctx context.Context, store *memory.Store, + id desire.Identity, owner string, +) { + t.Helper() + PutReadDesire(t, ctx, store, id, owner, 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.TB, ctx context.Context, store *memory.Store, + id desire.Identity, owner string, content []byte, +) { + t.Helper() + PutReadDesire(t, ctx, store, id, owner, desire.ReadStatus{ + Status: desire.Status{Conditions: []metav1.Condition{{ + Type: desire.TypeSuccessful, Status: metav1.ConditionFalse, Reason: desire.ReasonKubeAPIError, + }}}, + KubeContent: content, + }) +} + +// PutDeleteDesire creates a delete desire with the given condition status and reason. +func PutDeleteDesire( + t testing.TB, ctx context.Context, store *memory.Store, + id desire.Identity, owner string, + condStatus metav1.ConditionStatus, reason string, +) { + t.Helper() + dd, err := store.CreateDeleteDesire(ctx, desire.DeleteDesire{Identity: id, Owner: owner}) + require.NoError(t, err) + _, err = store.UpdateDeleteDesireStatus(ctx, id, desire.Status{ + Conditions: []metav1.Condition{{ + Type: desire.TypeSuccessful, Status: condStatus, Reason: reason, + }}, + }, dd.Version) + require.NoError(t, err) +} + +// PutConfirmedDeleteDesire creates a delete desire marked as successfully deleted. +func PutConfirmedDeleteDesire( + t testing.TB, ctx context.Context, store *memory.Store, + id desire.Identity, owner string, +) { + t.Helper() + PutDeleteDesire(t, ctx, store, id, owner, metav1.ConditionTrue, desire.ReasonDeleted) +} + +// MarkReadDesireNotFound updates an existing read desire to NotFound status. +func MarkReadDesireNotFound( + t testing.TB, ctx context.Context, store *memory.Store, id desire.Identity, +) { + t.Helper() + _, err := store.UpdateReadDesireStatus(ctx, id, desire.ReadStatus{ + Status: desire.Status{Conditions: []metav1.Condition{{ + Type: desire.TypeSuccessful, Status: metav1.ConditionFalse, Reason: desire.ReasonNotFound, + }}}, + }) + require.NoError(t, err) +} + +// MarkReadDesireSynced updates an existing read desire to Synced status with content. +func MarkReadDesireSynced( + t testing.TB, ctx context.Context, store *memory.Store, + id desire.Identity, content []byte, +) { + t.Helper() + _, err := store.UpdateReadDesireStatus(ctx, id, desire.ReadStatus{ + Status: desire.Status{Conditions: []metav1.Condition{{ + Type: desire.TypeSuccessful, Status: metav1.ConditionTrue, Reason: desire.ReasonSynced, + }}}, + KubeContent: content, + }) + require.NoError(t, err) +} + +// MarkDeleteDesireConfirmed updates an existing delete desire to confirmed-deleted status. +func MarkDeleteDesireConfirmed( + t testing.TB, ctx context.Context, store *memory.Store, id desire.Identity, +) { + t.Helper() + dd, err := store.GetDeleteDesire(ctx, id) + require.NoError(t, err) + _, err = store.UpdateDeleteDesireStatus(ctx, id, desire.Status{ + Conditions: []metav1.Condition{{ + Type: desire.TypeSuccessful, Status: metav1.ConditionTrue, Reason: desire.ReasonDeleted, + }}, + }, dd.Version) + require.NoError(t, err) +} + +// 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} +} + +// TestIdentity is a builder for desire.Identity values in tests. +// The With* methods return copies, so the original is never mutated. +type TestIdentity struct { + ManagementCluster string + Resource string + Namespace string + Name string +} + +func (ti TestIdentity) build(t desire.DesireType) desire.Identity { + return desire.Identity{ + ManagementCluster: ti.ManagementCluster, Type: t, + Resource: ti.Resource, Namespace: ti.Namespace, Name: ti.Name, + } +} + +func (ti TestIdentity) Read() desire.Identity { return ti.build(desire.TypeRead) } +func (ti TestIdentity) Delete() desire.Identity { return ti.build(desire.TypeDelete) } +func (ti TestIdentity) Apply() desire.Identity { return ti.build(desire.TypeApply) } + +func (ti TestIdentity) WithName(name string) TestIdentity { ti.Name = name; return ti } +func (ti TestIdentity) WithNamespace(namespace string) TestIdentity { ti.Namespace = namespace; return ti } From 2f67f45dea39772c6b741f1aba21ff12819a0483 Mon Sep 17 00:00:00 2001 From: Michal Vavrinec Date: Fri, 18 Sep 2026 14:02:20 +0200 Subject: [PATCH 3/8] HYPERFLEET-1439 - feat: add ErrDeletionPending to desire cleanup lifecycle MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Guard cleanup against orphaning resources when the apply desire still exists and the applier may not have processed it yet. Treat pending deletion as a transient state — log warn instead of error and skip the deletion error metric. Co-Authored-By: Claude Opus 4.6 --- internal/desireclient/cleanup.go | 24 ++++++++++++++++++++---- internal/desireclient/cleanup_test.go | 21 +++++++++++++++++++++ internal/desireclient/types.go | 7 +++++++ internal/executor/resource_executor.go | 11 ++++++++--- 4 files changed, 56 insertions(+), 7 deletions(-) diff --git a/internal/desireclient/cleanup.go b/internal/desireclient/cleanup.go index c5bc472d..a31b9c2b 100644 --- a/internal/desireclient/cleanup.go +++ b/internal/desireclient/cleanup.go @@ -14,7 +14,9 @@ import ( // CleanupAfterDeletion implements transportclient.DesireCleaner. It removes // the delete desire (only when the applier confirms deletion) then the read // desire. Returns an error if the delete desire exists but is not yet confirmed, -// causing the executor to retry on the next reconciliation. +// or if no delete desire exists but an apply desire is still present (the +// applier may not have applied it yet), causing the executor to retry on the +// next reconciliation. func (c *Client) CleanupAfterDeletion( ctx context.Context, gvk schema.GroupVersionKind, @@ -34,13 +36,27 @@ func (c *Client) CleanupAfterDeletion( dd, err := c.store.GetDeleteDesire(ctx, deleteID) switch { case errors.Is(err, desire.ErrNotFound): - // No delete desire — proceed to read desire cleanup. + applyID, buildErr := buildIdentity(tc, desire.TypeApply, gvk, namespace, name) + if buildErr != nil { + return buildErr + } + _, applyErr := c.store.GetApplyDesire(ctx, applyID) + switch { + case applyErr == nil: + return fmt.Errorf( + "desireclient: cleanup: apply desire still exists for %s/%s,"+ + " resource may not have been created yet: %w", + namespace, name, ErrDeletionPending) + case !errors.Is(applyErr, desire.ErrNotFound): + return fmt.Errorf("desireclient: cleanup: failed to get apply desire for %s/%s: %w", + namespace, name, applyErr) + } case err != nil: return fmt.Errorf("desireclient: cleanup: failed to get delete desire for %s/%s: %w", namespace, name, err) case !desire.IsDeleted(dd.Status): - return fmt.Errorf("desireclient: cleanup: deletion not yet confirmed for %s/%s", - namespace, name) + return fmt.Errorf("desireclient: cleanup: deletion not yet confirmed for %s/%s: %w", + namespace, name, ErrDeletionPending) default: if delErr := c.store.DeleteDeleteDesire(ctx, deleteID, c.owner, dd.Version); delErr != nil { return fmt.Errorf("desireclient: cleanup: failed to delete delete desire for %s/%s: %w", diff --git a/internal/desireclient/cleanup_test.go b/internal/desireclient/cleanup_test.go index 0ff9bce3..cb721cce 100644 --- a/internal/desireclient/cleanup_test.go +++ b/internal/desireclient/cleanup_test.go @@ -213,6 +213,17 @@ func TestCleanupAfterDeletion_DeleteReadDesireError(t *testing.T) { assert.Contains(t, err.Error(), "failed to delete read desire") } +func TestCleanupAfterDeletion_GetApplyDesireError_NoDeleteDesire_ReturnsStoreError(t *testing.T) { + ctx := context.Background() + store := &failingGetApplyDesireStore{SpecStore: newMemoryStore()} + c := newTestClient(store) + + err := c.CleanupAfterDeletion(ctx, testGVK(), testNamespace, testName, testTransportContext()) + require.Error(t, err) + assert.Contains(t, err.Error(), "failed to get apply desire") + assert.False(t, errors.Is(err, ErrDeletionPending), "genuine store error must not be wrapped as ErrDeletionPending") +} + // --- Test store wrappers --- type failingGetDeleteDesireStore struct { @@ -225,6 +236,16 @@ func (f *failingGetDeleteDesireStore) GetDeleteDesire( return desire.DeleteDesire{}, errors.New("boom: store unavailable") } +type failingGetApplyDesireStore struct { + desire.SpecStore +} + +func (f *failingGetApplyDesireStore) GetApplyDesire( + _ context.Context, _ desire.Identity, +) (desire.ApplyDesire, error) { + return desire.ApplyDesire{}, errors.New("boom: store unavailable") +} + type failingDeleteDeleteDesireStore struct { desire.SpecStore } diff --git a/internal/desireclient/types.go b/internal/desireclient/types.go index e8f9bfe1..75962c41 100644 --- a/internal/desireclient/types.go +++ b/internal/desireclient/types.go @@ -37,6 +37,13 @@ type TransportContext struct { // errors.Is to check for it. var ErrNotSyncedYet = errors.New("desireclient: resource not synced yet") +// ErrDeletionPending indicates that the desire transport has not yet +// confirmed resource deletion — either the apply desire still exists +// (applier may not have processed it) or the delete desire exists but +// the applier has not confirmed deletion. This is an expected transient +// state during the deletion lifecycle, not a hard failure. +var ErrDeletionPending = errors.New("desireclient: deletion pending") + // resolveTransportContext type-asserts the generic TransportContext and // validates both fields are set. func resolveTransportContext(target transportclient.TransportContext) (*TransportContext, error) { diff --git a/internal/executor/resource_executor.go b/internal/executor/resource_executor.go index 144de7e0..fa743daa 100644 --- a/internal/executor/resource_executor.go +++ b/internal/executor/resource_executor.go @@ -721,12 +721,17 @@ func (re *ResourceExecutor) executeResourceDelete( execCtx.Resources[resource.Name] = nil if err := re.tryCleanupDesires(ctx, resource, execCtx, transportClient, transportTarget, gvk); err != nil { - slog.ErrorContext(ctx, "resource desire cleanup failed after delete", - "resource", resource.Name, "error", err) + if errors.Is(err, desireclient.ErrDeletionPending) { + slog.WarnContext(ctx, "resource desire cleanup: deletion pending", + "resource", resource.Name, "error", err) + } else { + slog.ErrorContext(ctx, "resource desire cleanup failed after delete", + "resource", resource.Name, "error", err) + re.metrics.RecordDeletion(resourceType, metrics.DeletionStatusError) + } result.Status = StatusFailed result.Error = err re.recordResourceError(execCtx, resource, err) - re.metrics.RecordDeletion(resourceType, metrics.DeletionStatusError) re.metrics.ObserveDeletionDuration(resourceType, time.Since(startTime)) return result, NewExecutorError(PhaseResources, resource.Name, "desire cleanup failed", err) } From ea8e5583030984c839fa76653b1df3bafff6ec59 Mon Sep 17 00:00:00 2001 From: Michal Vavrinec Date: Fri, 18 Sep 2026 14:04:49 +0200 Subject: [PATCH 4/8] HYPERFLEET-1675 - feat: handle ErrNotSyncedYet across executor lifecycle Skip post-apply discovery and treat as absent in pre-discovery when the applier hasn't synced the read mirror yet. Log warn and skip error metric in the delete path. - Add PutUnsyncedReadDesire test helper - Add 4 desire transport lifecycle tests including full empty-store cycle Co-Authored-By: Claude Opus 4.6 --- internal/desireclient/discover_test.go | 6 +- internal/desireclient/testhelpers.go | 24 +- internal/executor/resource_executor.go | 14 +- internal/executor/resource_executor_test.go | 335 ++++++++++++++++++++ 4 files changed, 371 insertions(+), 8 deletions(-) diff --git a/internal/desireclient/discover_test.go b/internal/desireclient/discover_test.go index add9d75b..933f7cad 100644 --- a/internal/desireclient/discover_test.go +++ b/internal/desireclient/discover_test.go @@ -59,8 +59,10 @@ func TestDiscoverResources_LabelSelectorMatchesSubset(t *testing.T) { "apiVersion": "v1", "kind": "ConfigMap", "metadata": {"name": "unlabeled", "namespace": "default"} }`) - PutSyncedReadDesire(t, ctx, store, testID.WithNamespace("labeled").WithName("labeled").Read(), testOwner, labeledManifest) - PutSyncedReadDesire(t, ctx, store, testID.WithNamespace("unlabeled").WithName("unlabeled").Read(), testOwner, unlabeledManifest) + PutSyncedReadDesire(t, ctx, store, + testID.WithNamespace("labeled").WithName("labeled").Read(), testOwner, labeledManifest) + PutSyncedReadDesire(t, ctx, store, + testID.WithNamespace("unlabeled").WithName("unlabeled").Read(), testOwner, unlabeledManifest) list, err := c.DiscoverResources(ctx, testGVK(), &manifest.DiscoveryConfig{LabelSelector: "app=myapp"}, testTransportContext()) diff --git a/internal/desireclient/testhelpers.go b/internal/desireclient/testhelpers.go index 044c6a0d..1689588f 100644 --- a/internal/desireclient/testhelpers.go +++ b/internal/desireclient/testhelpers.go @@ -10,6 +10,19 @@ import ( metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" ) +// PutUnsyncedReadDesire creates a read desire with no status condition, +// simulating a desire the applier has not observed yet. +func PutUnsyncedReadDesire( + t testing.TB, ctx context.Context, store *memory.Store, + id desire.Identity, owner string, +) { + t.Helper() + _, err := store.CreateReadDesire(ctx, desire.ReadDesire{ + Identity: id, Owner: owner, TargetVersion: "v1", + }) + require.NoError(t, err) +} + // PutReadDesire creates a read desire with the given identity, owner, and status. func PutReadDesire( t testing.TB, ctx context.Context, store *memory.Store, @@ -178,8 +191,11 @@ func (ti TestIdentity) build(t desire.DesireType) desire.Identity { } func (ti TestIdentity) Read() desire.Identity { return ti.build(desire.TypeRead) } -func (ti TestIdentity) Delete() desire.Identity { return ti.build(desire.TypeDelete) } -func (ti TestIdentity) Apply() desire.Identity { return ti.build(desire.TypeApply) } +func (ti TestIdentity) Delete() desire.Identity { return ti.build(desire.TypeDelete) } +func (ti TestIdentity) Apply() desire.Identity { return ti.build(desire.TypeApply) } -func (ti TestIdentity) WithName(name string) TestIdentity { ti.Name = name; return ti } -func (ti TestIdentity) WithNamespace(namespace string) TestIdentity { ti.Namespace = namespace; return ti } +func (ti TestIdentity) WithName(name string) TestIdentity { ti.Name = name; return ti } +func (ti TestIdentity) WithNamespace(namespace string) TestIdentity { + ti.Namespace = namespace + return ti +} diff --git a/internal/executor/resource_executor.go b/internal/executor/resource_executor.go index fa743daa..261d750e 100644 --- a/internal/executor/resource_executor.go +++ b/internal/executor/resource_executor.go @@ -219,6 +219,11 @@ func (re *ResourceExecutor) executeResource( // Step 7: Post-apply discovery — find the applied resource and store in execCtx for CEL evaluation if resource.Discovery != nil { discovered, discoverErr := re.discoverResource(ctx, resource, execCtx, transportClient, transportTarget) + if errors.Is(discoverErr, desireclient.ErrNotSyncedYet) { + slog.DebugContext(ctx, "resource not synced yet, skipping post-apply discovery", + "resource", resource.Name) + return result, nil + } if discoverErr != nil { result.Status = StatusFailed result.Error = discoverErr @@ -560,7 +565,7 @@ func (re *ResourceExecutor) preDiscoverAll( discovered, err := re.discoverResource(ctx, resource, execCtx, transportClient, transportTarget) if err != nil { - if apierrors.IsNotFound(err) { + if apierrors.IsNotFound(err) || errors.Is(err, desireclient.ErrNotSyncedYet) { // Resource does not exist yet — leave absent from context. continue } @@ -705,10 +710,15 @@ func (re *ResourceExecutor) executeResourceDelete( isNotFound := discoverErr != nil && apierrors.IsNotFound(discoverErr) if discoverErr != nil && !isNotFound { + if errors.Is(discoverErr, desireclient.ErrNotSyncedYet) { + slog.WarnContext(ctx, "resource not synced yet, cannot discover for deletion", + "resource", resource.Name, "error", discoverErr) + } else { + re.metrics.RecordDeletion(resourceType, metrics.DeletionStatusError) + } result.Status = StatusFailed result.Error = discoverErr re.recordResourceError(execCtx, resource, discoverErr) - re.metrics.RecordDeletion(resourceType, metrics.DeletionStatusError) re.metrics.ObserveDeletionDuration(resourceType, time.Since(startTime)) return result, NewExecutorError( PhaseResources, resource.Name, "failed to discover resource for deletion", discoverErr) diff --git a/internal/executor/resource_executor_test.go b/internal/executor/resource_executor_test.go index ee593861..3b3358f4 100644 --- a/internal/executor/resource_executor_test.go +++ b/internal/executor/resource_executor_test.go @@ -10,9 +10,12 @@ import ( "github.com/openshift-hyperfleet/hyperfleet-adapter/internal/k8sclient" "github.com/openshift-hyperfleet/hyperfleet-adapter/internal/manifest" "github.com/openshift-hyperfleet/hyperfleet-adapter/internal/transportclient" + "github.com/openshift-hyperfleet/hyperfleet-applier/pkg/desire" + "github.com/openshift-hyperfleet/hyperfleet-applier/pkg/desire/store/memory" "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" "k8s.io/apimachinery/pkg/apis/meta/v1/unstructured" "k8s.io/apimachinery/pkg/runtime/schema" ) @@ -2452,3 +2455,335 @@ func TestResourceExecutor_LifecycleDelete_StillPresent_NoCleanup(t *testing.T) { assert.Equal(t, StatusSuccess, results[0].Status) assert.False(t, mock.CleanupCalled, "CleanupAfterDeletion must not be called when resource is still present") } + +// ---- Desire transport lifecycle tests ---- + +const ( + desireTransportName = "desire-primary" + desireOwner = "hyperfleet-adapter" +) + +var testDesireID = desireclient.TestIdentity{ + ManagementCluster: "cluster-1", + Resource: "configmaps", + Namespace: "default", + Name: "test-config", +} + +func configMapContent() []byte { + return []byte(`{ + "apiVersion": "v1", + "kind": "ConfigMap", + "metadata": { + "name": "test-config", + "namespace": "default", + "annotations": {"hyperfleet.io/generation": "1"} + }, + "data": {"key": "value"} + }`) +} + +func newDesireExecutor(store *memory.Store) *ResourceExecutor { + c := desireclient.NewClient(store, desireOwner) + return newResourceExecutor(&ExecutorConfig{ + Config: &configloader.Config{ + Transports: map[string]configloader.TransportDefinition{ + desireTransportName: {Type: configloader.TransportTypeRemote}, + }, + }, + TransportRegistry: transportclient.Registry{ + desireTransportName: c, + }, + }) +} + +func newDesireResourceWithLifecycle(expression, propagationPolicy string) configloader.Resource { + r := configloader.Resource{ + Name: "test-resource", + Transport: &configloader.TransportConfig{ + Client: desireTransportName, + Desire: &configloader.DesireTransportConfig{ + TargetCluster: testDesireID.ManagementCluster, + Resource: testDesireID.Resource, + }, + }, + Manifest: map[string]interface{}{ + "apiVersion": "v1", + "kind": "ConfigMap", + "metadata": map[string]interface{}{ + "name": testDesireID.Name, + "namespace": testDesireID.Namespace, + "annotations": map[string]interface{}{ + "hyperfleet.io/generation": "1", + }, + }, + "data": map[string]interface{}{"key": "value"}, + }, + Discovery: &configloader.DiscoveryConfig{ + Namespace: testDesireID.Namespace, + ByName: testDesireID.Name, + }, + Lifecycle: &configloader.ResourceLifecycle{ + Delete: &configloader.LifecycleDelete{ + PropagationPolicy: propagationPolicy, + }, + }, + } + if expression != "" { + r.Lifecycle.Delete.When = &configloader.LifecycleWhen{Expression: expression} + } + return r +} + +// Test 1: Slow applier — three events to complete the full apply → delete → cleanup cycle. +func TestResourceExecutor_DesireTransport_SlowApplier(t *testing.T) { + ctx := context.Background() + store := memory.New() + re := newDesireExecutor(store) + resource := newDesireResourceWithLifecycle("deleted_time != null", "Background") + + desireclient.PutUnsyncedReadDesire(t, ctx, store, testDesireID.Read(), desireOwner) + + // ---- Event 1: apply (deleted_time absent → delete.when false) ---- + execCtx := NewExecutionContext(ctx, nil, nil) + results, err := re.ExecuteAll(ctx, []configloader.Resource{resource}, execCtx) + require.NoError(t, err) + require.Len(t, results, 1) + assert.Equal(t, StatusSuccess, results[0].Status) + + _, err = store.GetApplyDesire(ctx, testDesireID.Apply()) + require.NoError(t, err, "ApplyDesire must exist after apply") + _, err = store.GetReadDesire(ctx, testDesireID.Read()) + require.NoError(t, err, "ReadDesire must exist after apply") + + desireclient.MarkReadDesireSynced(t, ctx, store, testDesireID.Read(), configMapContent()) + + // ---- Event 2: delete (applier hasn't confirmed yet) ---- + execCtx = NewExecutionContext(ctx, nil, nil) + execCtx.Params["deleted_time"] = testDeletedTime + results, err = re.ExecuteAll(ctx, []configloader.Resource{resource}, execCtx) + require.NoError(t, err) + require.Len(t, results, 1) + assert.Equal(t, StatusSuccess, results[0].Status) + + _, err = store.GetApplyDesire(ctx, testDesireID.Apply()) + assert.ErrorIs(t, err, desire.ErrNotFound, "ApplyDesire must be removed by DeleteResource") + + _, err = store.GetDeleteDesire(ctx, testDesireID.Delete()) + assert.NoError(t, err, "DeleteDesire must exist (pending)") + + _, err = store.GetReadDesire(ctx, testDesireID.Read()) + assert.NoError(t, err, "ReadDesire must still exist") + + desireclient.MarkDeleteDesireConfirmed(t, ctx, store, testDesireID.Delete()) + desireclient.MarkReadDesireNotFound(t, ctx, store, testDesireID.Read()) + + // ---- Event 3: delete (applier confirmed) ---- + execCtx = NewExecutionContext(ctx, nil, nil) + execCtx.Params["deleted_time"] = testDeletedTime + results, err = re.ExecuteAll(ctx, []configloader.Resource{resource}, execCtx) + require.NoError(t, err) + require.Len(t, results, 1) + assert.Equal(t, StatusSuccess, results[0].Status) + assert.Equal(t, "resource already deleted or never existed", results[0].OperationReason) + + _, err = store.GetDeleteDesire(ctx, testDesireID.Delete()) + assert.ErrorIs(t, err, desire.ErrNotFound, "DeleteDesire must be removed by cleanup") + + _, err = store.GetReadDesire(ctx, testDesireID.Read()) + assert.ErrorIs(t, err, desire.ErrNotFound, "ReadDesire must be removed by cleanup") +} + +// instantApplierStore wraps memory.Store. When a DeleteDesire is created it +// immediately marks it Deleted and the paired ReadDesire as NotFound, +// simulating an applier that confirms before post-delete discovery runs. +type instantApplierStore struct { + *memory.Store +} + +func (s *instantApplierStore) CreateDeleteDesire( + ctx context.Context, dd desire.DeleteDesire, +) (desire.DeleteDesire, error) { + created, err := s.Store.CreateDeleteDesire(ctx, dd) + if err != nil { + return created, err + } + s.UpdateDeleteDesireStatus(ctx, dd.Identity, desire.Status{ + Conditions: []metav1.Condition{{ + Type: desire.TypeSuccessful, Status: metav1.ConditionTrue, Reason: desire.ReasonDeleted, + }}, + }, created.Version) + + readID := dd.Identity + readID.Type = desire.TypeRead + s.UpdateReadDesireStatus(ctx, readID, desire.ReadStatus{ + Status: desire.Status{Conditions: []metav1.Condition{{ + Type: desire.TypeSuccessful, Status: metav1.ConditionFalse, Reason: desire.ReasonNotFound, + }}}, + }) + return created, nil +} + +// Test 2: Fast applier — applier confirms between DeleteResource and +// post-delete discovery within the same event. +func TestResourceExecutor_DesireTransport_FastApplier(t *testing.T) { + ctx := context.Background() + inner := memory.New() + store := &instantApplierStore{Store: inner} + c := desireclient.NewClient(store, desireOwner) + re := newResourceExecutor(&ExecutorConfig{ + Config: &configloader.Config{ + Transports: map[string]configloader.TransportDefinition{ + desireTransportName: {Type: configloader.TransportTypeRemote}, + }, + }, + TransportRegistry: transportclient.Registry{ + desireTransportName: c, + }, + }) + resource := newDesireResourceWithLifecycle("deleted_time != null", "Background") + + desireclient.PutUnsyncedReadDesire(t, ctx, inner, testDesireID.Read(), desireOwner) + + // ---- Event 1: apply ---- + execCtx := NewExecutionContext(ctx, nil, nil) + results, err := re.ExecuteAll(ctx, []configloader.Resource{resource}, execCtx) + require.NoError(t, err) + require.Len(t, results, 1) + assert.Equal(t, StatusSuccess, results[0].Status) + + desireclient.MarkReadDesireSynced(t, ctx, inner, testDesireID.Read(), configMapContent()) + + // ---- Event 2: delete (applier confirms instantly via wrapper) ---- + execCtx = NewExecutionContext(ctx, nil, nil) + execCtx.Params["deleted_time"] = testDeletedTime + results, err = re.ExecuteAll(ctx, []configloader.Resource{resource}, execCtx) + require.NoError(t, err) + require.Len(t, results, 1) + assert.Equal(t, StatusSuccess, results[0].Status) + + _, err = inner.GetDeleteDesire(ctx, testDesireID.Delete()) + assert.ErrorIs(t, err, desire.ErrNotFound, "DeleteDesire must be removed in single cycle") + + _, err = inner.GetReadDesire(ctx, testDesireID.Read()) + assert.ErrorIs(t, err, desire.ErrNotFound, "ReadDesire must be removed in single cycle") +} + +// Test 3: Transient NotFound before apply lands — regression test. +// +// The applier's read informer ran before the apply pass and wrote +// Reason=NotFound on the ReadDesire. The ApplyDesire is still live. +// Cleanup must refuse because the apply desire exists, preventing +// an orphaned resource on the target cluster. +func TestResourceExecutor_DesireTransport_TransientNotFound(t *testing.T) { + ctx := context.Background() + store := memory.New() + re := newDesireExecutor(store) + resource := newDesireResourceWithLifecycle("deleted_time != null", "Background") + + desireclient.PutUnsyncedReadDesire(t, ctx, store, testDesireID.Read(), desireOwner) + + // ---- Event 1: apply ---- + execCtx := NewExecutionContext(ctx, nil, nil) + results, err := re.ExecuteAll(ctx, []configloader.Resource{resource}, execCtx) + require.NoError(t, err) + require.Len(t, results, 1) + assert.Equal(t, StatusSuccess, results[0].Status) + + // Applier's read informer ran before the apply pass: ReadDesire + // stays NotFound. ApplyDesire exists but hasn't been applied yet. + + _, err = store.GetApplyDesire(ctx, testDesireID.Apply()) + require.NoError(t, err, "ApplyDesire must exist after apply") + + // ---- Event 2: delete.when true ---- + // Discovery reads mirror → NotFound → step 2 → tryCleanupDesires. + // CleanupAfterDeletion sees no DeleteDesire but ApplyDesire exists → error. + execCtx = NewExecutionContext(ctx, nil, nil) + execCtx.Params["deleted_time"] = testDeletedTime + results, err = re.ExecuteAll(ctx, []configloader.Resource{resource}, execCtx) + require.Error(t, err) + require.Len(t, results, 1) + assert.Equal(t, StatusFailed, results[0].Status) + + _, err = store.GetApplyDesire(ctx, testDesireID.Apply()) + assert.NoError(t, err, "ApplyDesire must survive — no orphan") + + _, err = store.GetReadDesire(ctx, testDesireID.Read()) + assert.NoError(t, err, "ReadDesire must survive — no orphan") +} + +// Test 4: Full lifecycle from an empty store — covers fresh-cluster apply +// and post-cleanup re-apply without seeding any read desire. +func TestResourceExecutor_DesireTransport_FullLifecycleFromEmptyStore(t *testing.T) { + ctx := context.Background() + store := memory.New() + re := newDesireExecutor(store) + resource := newDesireResourceWithLifecycle("deleted_time != null", "Background") + + // ---- Event 1: apply on empty store (no read desire exists) ---- + execCtx := NewExecutionContext(ctx, nil, nil) + results, err := re.ExecuteAll(ctx, []configloader.Resource{resource}, execCtx) + require.NoError(t, err) + require.Len(t, results, 1) + assert.Equal(t, StatusSuccess, results[0].Status) + + _, err = store.GetApplyDesire(ctx, testDesireID.Apply()) + require.NoError(t, err, "ApplyDesire must exist") + _, err = store.GetReadDesire(ctx, testDesireID.Read()) + require.NoError(t, err, "ReadDesire must be auto-created by ensureReadDesire") + + // Simulate applier syncing the read mirror. + desireclient.MarkReadDesireSynced(t, ctx, store, testDesireID.Read(), configMapContent()) + + // ---- Event 2: apply after sync (post-apply discovery finds the resource) ---- + execCtx = NewExecutionContext(ctx, nil, nil) + results, err = re.ExecuteAll(ctx, []configloader.Resource{resource}, execCtx) + require.NoError(t, err) + require.Len(t, results, 1) + assert.Equal(t, StatusSuccess, results[0].Status) + assert.NotNil(t, execCtx.Resources["test-resource"], "synced resource must be stored in context") + + // ---- Event 3: delete (applier hasn't confirmed yet) ---- + execCtx = NewExecutionContext(ctx, nil, nil) + execCtx.Params["deleted_time"] = testDeletedTime + results, err = re.ExecuteAll(ctx, []configloader.Resource{resource}, execCtx) + require.NoError(t, err) + require.Len(t, results, 1) + assert.Equal(t, StatusSuccess, results[0].Status) + + _, err = store.GetApplyDesire(ctx, testDesireID.Apply()) + assert.ErrorIs(t, err, desire.ErrNotFound, "ApplyDesire must be removed") + _, err = store.GetDeleteDesire(ctx, testDesireID.Delete()) + assert.NoError(t, err, "DeleteDesire must exist (pending)") + + // Simulate applier confirming deletion. + desireclient.MarkDeleteDesireConfirmed(t, ctx, store, testDesireID.Delete()) + desireclient.MarkReadDesireNotFound(t, ctx, store, testDesireID.Read()) + + // ---- Event 4: delete (applier confirmed → cleanup removes all desires) ---- + execCtx = NewExecutionContext(ctx, nil, nil) + execCtx.Params["deleted_time"] = testDeletedTime + results, err = re.ExecuteAll(ctx, []configloader.Resource{resource}, execCtx) + require.NoError(t, err) + require.Len(t, results, 1) + assert.Equal(t, StatusSuccess, results[0].Status) + assert.Equal(t, "resource already deleted or never existed", results[0].OperationReason) + + _, err = store.GetDeleteDesire(ctx, testDesireID.Delete()) + assert.ErrorIs(t, err, desire.ErrNotFound, "DeleteDesire must be removed by cleanup") + _, err = store.GetReadDesire(ctx, testDesireID.Read()) + assert.ErrorIs(t, err, desire.ErrNotFound, "ReadDesire must be removed by cleanup") + + // ---- Event 5: apply on empty store again (post-cleanup) ---- + execCtx = NewExecutionContext(ctx, nil, nil) + results, err = re.ExecuteAll(ctx, []configloader.Resource{resource}, execCtx) + require.NoError(t, err) + require.Len(t, results, 1) + assert.Equal(t, StatusSuccess, results[0].Status) + + _, err = store.GetApplyDesire(ctx, testDesireID.Apply()) + require.NoError(t, err, "ApplyDesire must exist after re-apply") + _, err = store.GetReadDesire(ctx, testDesireID.Read()) + require.NoError(t, err, "ReadDesire must be re-created after cleanup") +} From 63e5ab4252005f8c02f8a752ae085e0cdce33721 Mon Sep 17 00:00:00 2001 From: Michal Vavrinec Date: Fri, 18 Sep 2026 14:52:22 +0200 Subject: [PATCH 5/8] HYPERFLEET-1675 - refactor: move test helpers to desiretest package Test helpers in testhelpers.go compiled into the production binary since it was not a _test.go file. Move them to a dedicated internal/desireclient/desiretest package so testing and testify are only linked in test builds. Co-Authored-By: Claude Opus 4.6 --- internal/desireclient/cleanup_test.go | 10 +++--- internal/desireclient/desireclient_test.go | 3 +- .../desiretest.go} | 2 +- internal/desireclient/discover_test.go | 27 +++++++------- internal/desireclient/get_test.go | 21 +++++------ internal/executor/resource_executor_test.go | 35 +++++++++++-------- 6 files changed, 54 insertions(+), 44 deletions(-) rename internal/desireclient/{testhelpers.go => desiretest/desiretest.go} (99%) diff --git a/internal/desireclient/cleanup_test.go b/internal/desireclient/cleanup_test.go index cb721cce..89d21567 100644 --- a/internal/desireclient/cleanup_test.go +++ b/internal/desireclient/cleanup_test.go @@ -5,6 +5,7 @@ import ( "errors" "testing" + "github.com/openshift-hyperfleet/hyperfleet-adapter/internal/desireclient/desiretest" "github.com/openshift-hyperfleet/hyperfleet-applier/pkg/desire" "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" @@ -25,7 +26,7 @@ func TestCleanupAfterDeletion_ConfirmedDelete_RemovesBoth(t *testing.T) { Resource: testResource, Namespace: testNamespace, Name: testName, } - PutConfirmedDeleteDesire(t, ctx, store, testID.Delete(), testOwner) + desiretest.PutConfirmedDeleteDesire(t, ctx, store, testID.Delete(), testOwner) _, err := store.CreateReadDesire(ctx, desire.ReadDesire{ Identity: readID, Owner: testOwner, TargetVersion: "v1", @@ -56,7 +57,8 @@ func TestCleanupAfterDeletion_PendingDelete_SkipsCleanup(t *testing.T) { Resource: testResource, Namespace: testNamespace, Name: testName, } - PutDeleteDesire(t, ctx, store, testID.Delete(), testOwner, metav1.ConditionFalse, desire.ReasonWaitingForDeletion) + desiretest.PutDeleteDesire(t, ctx, store, testID.Delete(), testOwner, + metav1.ConditionFalse, desire.ReasonWaitingForDeletion) _, err := store.CreateReadDesire(ctx, desire.ReadDesire{ Identity: readID, Owner: testOwner, TargetVersion: "v1", @@ -151,7 +153,7 @@ func TestCleanupAfterDeletion_DeleteDesireOnly_NoReadDesire(t *testing.T) { Resource: testResource, Namespace: testNamespace, Name: testName, } - PutConfirmedDeleteDesire(t, ctx, store, testID.Delete(), testOwner) + desiretest.PutConfirmedDeleteDesire(t, ctx, store, testID.Delete(), testOwner) err := c.CleanupAfterDeletion(ctx, testGVK(), testNamespace, testName, testTransportContext()) require.NoError(t, err) @@ -182,7 +184,7 @@ func TestCleanupAfterDeletion_DeleteDeleteDesireError(t *testing.T) { ctx := context.Background() inner := newMemoryStore() - PutConfirmedDeleteDesire(t, ctx, inner, testID.Delete(), testOwner) + desiretest.PutConfirmedDeleteDesire(t, ctx, inner, testID.Delete(), testOwner) store := &failingDeleteDeleteDesireStore{SpecStore: inner} c := newTestClient(store) diff --git a/internal/desireclient/desireclient_test.go b/internal/desireclient/desireclient_test.go index b837037d..bb79e8eb 100644 --- a/internal/desireclient/desireclient_test.go +++ b/internal/desireclient/desireclient_test.go @@ -5,6 +5,7 @@ import ( "errors" "fmt" + "github.com/openshift-hyperfleet/hyperfleet-adapter/internal/desireclient/desiretest" "github.com/openshift-hyperfleet/hyperfleet-adapter/pkg/constants" "github.com/openshift-hyperfleet/hyperfleet-applier/pkg/desire" "github.com/openshift-hyperfleet/hyperfleet-applier/pkg/desire/store/memory" @@ -72,7 +73,7 @@ func (f *failingCreateDeleteDesireStore) CreateDeleteDesire( return desire.DeleteDesire{}, errors.New("boom: delete desire store unavailable") } -var testID = TestIdentity{ +var testID = desiretest.TestIdentity{ ManagementCluster: testManagementCluster, Resource: testResource, Namespace: testNamespace, diff --git a/internal/desireclient/testhelpers.go b/internal/desireclient/desiretest/desiretest.go similarity index 99% rename from internal/desireclient/testhelpers.go rename to internal/desireclient/desiretest/desiretest.go index 1689588f..242f723f 100644 --- a/internal/desireclient/testhelpers.go +++ b/internal/desireclient/desiretest/desiretest.go @@ -1,4 +1,4 @@ -package desireclient +package desiretest import ( "context" diff --git a/internal/desireclient/discover_test.go b/internal/desireclient/discover_test.go index 933f7cad..96927266 100644 --- a/internal/desireclient/discover_test.go +++ b/internal/desireclient/discover_test.go @@ -4,6 +4,7 @@ import ( "context" "testing" + "github.com/openshift-hyperfleet/hyperfleet-adapter/internal/desireclient/desiretest" "github.com/openshift-hyperfleet/hyperfleet-adapter/internal/manifest" "github.com/openshift-hyperfleet/hyperfleet-applier/pkg/desire" "github.com/stretchr/testify/assert" @@ -25,7 +26,7 @@ func TestDiscoverResources_ReturnsSyncedResourceByName(t *testing.T) { store := newMemoryStore() c := newTestClient(store) - PutSyncedReadDesire(t, ctx, store, testID.Read(), testOwner, configMapManifest(1)) + desiretest.PutSyncedReadDesire(t, ctx, store, testID.Read(), testOwner, configMapManifest(1)) list, err := c.DiscoverResources(ctx, testGVK(), &manifest.DiscoveryConfig{ByName: testName}, testTransportContext()) require.NoError(t, err) @@ -38,7 +39,7 @@ func TestDiscoverResources_ByNameExcludesNonMatchingName(t *testing.T) { store := newMemoryStore() c := newTestClient(store) - PutSyncedReadDesire(t, ctx, store, testID.Read(), testOwner, configMapManifest(1)) + desiretest.PutSyncedReadDesire(t, ctx, store, testID.Read(), testOwner, configMapManifest(1)) discovery := &manifest.DiscoveryConfig{ByName: "other-name"} list, err := c.DiscoverResources(ctx, testGVK(), discovery, testTransportContext()) @@ -59,9 +60,9 @@ func TestDiscoverResources_LabelSelectorMatchesSubset(t *testing.T) { "apiVersion": "v1", "kind": "ConfigMap", "metadata": {"name": "unlabeled", "namespace": "default"} }`) - PutSyncedReadDesire(t, ctx, store, + desiretest.PutSyncedReadDesire(t, ctx, store, testID.WithNamespace("labeled").WithName("labeled").Read(), testOwner, labeledManifest) - PutSyncedReadDesire(t, ctx, store, + desiretest.PutSyncedReadDesire(t, ctx, store, testID.WithNamespace("unlabeled").WithName("unlabeled").Read(), testOwner, unlabeledManifest) list, err := c.DiscoverResources(ctx, testGVK(), @@ -93,7 +94,7 @@ func TestDiscoverResources_SurfacesRetainedMirrorOnFailedRead(t *testing.T) { store := newMemoryStore() c := newTestClient(store) - PutKubeAPIErrorReadDesire(t, ctx, store, testID.Read(), testOwner, configMapManifest(1)) + desiretest.PutKubeAPIErrorReadDesire(t, ctx, store, testID.Read(), testOwner, configMapManifest(1)) list, err := c.DiscoverResources(ctx, testGVK(), &manifest.DiscoveryConfig{}, testTransportContext()) require.NoError(t, err) @@ -107,7 +108,7 @@ func TestDiscoverResources_SkipsFailedReadWithNoRetainedMirror(t *testing.T) { store := newMemoryStore() c := newTestClient(store) - PutKubeAPIErrorReadDesire(t, ctx, store, testID.Read(), testOwner, nil) + desiretest.PutKubeAPIErrorReadDesire(t, ctx, store, testID.Read(), testOwner, nil) list, err := c.DiscoverResources(ctx, testGVK(), &manifest.DiscoveryConfig{}, testTransportContext()) require.NoError(t, err) @@ -119,7 +120,7 @@ func TestDiscoverResources_SkipsNotFoundFalseDesire(t *testing.T) { store := newMemoryStore() c := newTestClient(store) - PutNotFoundReadDesire(t, ctx, store, testID.Read(), testOwner) + desiretest.PutNotFoundReadDesire(t, ctx, store, testID.Read(), testOwner) list, err := c.DiscoverResources(ctx, testGVK(), &manifest.DiscoveryConfig{}, testTransportContext()) require.NoError(t, err) @@ -131,8 +132,8 @@ func TestDiscoverResources_SkipsUndecodableContentButKeepsOthers(t *testing.T) { store := newMemoryStore() c := newTestClient(store) - PutSyncedReadDesire(t, ctx, store, testID.WithName("bad").Read(), testOwner, []byte("not-json")) - PutSyncedReadDesire(t, ctx, store, testID.WithName("good").Read(), testOwner, configMapManifest(1)) + desiretest.PutSyncedReadDesire(t, ctx, store, testID.WithName("bad").Read(), testOwner, []byte("not-json")) + desiretest.PutSyncedReadDesire(t, ctx, store, testID.WithName("good").Read(), testOwner, 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") @@ -145,7 +146,7 @@ func TestDiscoverResources_FiltersOutOtherResourceType(t *testing.T) { store := newMemoryStore() c := newTestClient(store) - PutSyncedReadDesire(t, ctx, store, testID.Read(), testOwner, configMapManifest(1)) + desiretest.PutSyncedReadDesire(t, ctx, store, testID.Read(), testOwner, configMapManifest(1)) otherContext := &TransportContext{ManagementCluster: testManagementCluster, Resource: "secrets"} list, err := c.DiscoverResources(ctx, testGVK(), &manifest.DiscoveryConfig{}, otherContext) @@ -184,7 +185,7 @@ func TestDiscoverResources_ScopedToPartition(t *testing.T) { store := newMemoryStore() c := newTestClient(store) - PutSyncedReadDesire(t, ctx, store, testID.Read(), testOwner, configMapManifest(1)) + desiretest.PutSyncedReadDesire(t, ctx, store, testID.Read(), testOwner, configMapManifest(1)) otherPartition := &TransportContext{ManagementCluster: "other-cluster", Resource: testResource} list, err := c.DiscoverResources(ctx, testGVK(), &manifest.DiscoveryConfig{}, otherPartition) @@ -199,8 +200,8 @@ func TestDiscoverResources_MultipleMatches(t *testing.T) { 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, testID.WithName("first").Read(), testOwner, first) - PutSyncedReadDesire(t, ctx, store, testID.WithName("second").Read(), testOwner, second) + desiretest.PutSyncedReadDesire(t, ctx, store, testID.WithName("first").Read(), testOwner, first) + desiretest.PutSyncedReadDesire(t, ctx, store, testID.WithName("second").Read(), testOwner, second) list, err := c.DiscoverResources(ctx, testGVK(), &manifest.DiscoveryConfig{}, testTransportContext()) require.NoError(t, err) diff --git a/internal/desireclient/get_test.go b/internal/desireclient/get_test.go index bbc3fa5f..14929c34 100644 --- a/internal/desireclient/get_test.go +++ b/internal/desireclient/get_test.go @@ -5,6 +5,7 @@ import ( "errors" "testing" + "github.com/openshift-hyperfleet/hyperfleet-adapter/internal/desireclient/desiretest" "github.com/openshift-hyperfleet/hyperfleet-applier/pkg/desire" "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" @@ -42,7 +43,7 @@ func TestGetResource_SyncedReturnsMirroredObject(t *testing.T) { store := newMemoryStore() c := newTestClient(store) - PutSyncedReadDesire(t, ctx, store, testID.Read(), testOwner, configMapManifest(1)) + desiretest.PutSyncedReadDesire(t, ctx, store, testID.Read(), testOwner, configMapManifest(1)) obj, err := c.GetResource(ctx, testGVK(), testNamespace, testName, testTransportContext()) require.NoError(t, err) @@ -55,7 +56,7 @@ func TestGetResource_ConfirmedNotFound(t *testing.T) { store := newMemoryStore() c := newTestClient(store) - PutConfirmedAbsentReadDesire(t, ctx, store, testID.Read(), testOwner) + desiretest.PutConfirmedAbsentReadDesire(t, ctx, store, testID.Read(), testOwner) _, err := c.GetResource(ctx, testGVK(), testNamespace, testName, testTransportContext()) require.Error(t, err) @@ -68,7 +69,7 @@ func TestGetResource_InvalidReadDesire(t *testing.T) { store := newMemoryStore() c := newTestClient(store) - PutInvalidReadDesire(t, ctx, store, testID.Read(), testOwner) + desiretest.PutInvalidReadDesire(t, ctx, store, testID.Read(), testOwner) _, err := c.GetResource(ctx, testGVK(), testNamespace, testName, testTransportContext()) require.Error(t, err) @@ -83,7 +84,7 @@ func TestGetResource_K8sAPIErrorWithRetainedMirrorReturnsStaleContent(t *testing store := newMemoryStore() c := newTestClient(store) - PutKubeAPIErrorReadDesire(t, ctx, store, testID.Read(), testOwner, configMapManifest(1)) + desiretest.PutKubeAPIErrorReadDesire(t, ctx, store, testID.Read(), testOwner, configMapManifest(1)) obj, err := c.GetResource(ctx, testGVK(), testNamespace, testName, testTransportContext()) require.NoError(t, err, @@ -97,7 +98,7 @@ func TestGetResource_K8sAPIErrorWithNoMirrorYetIsNotSyncedYet(t *testing.T) { store := newMemoryStore() c := newTestClient(store) - PutKubeAPIErrorReadDesire(t, ctx, store, testID.Read(), testOwner, nil) + desiretest.PutKubeAPIErrorReadDesire(t, ctx, store, testID.Read(), testOwner, nil) _, err := c.GetResource(ctx, testGVK(), testNamespace, testName, testTransportContext()) require.Error(t, err) @@ -152,30 +153,30 @@ func TestDecodeReadDesire(t *testing.T) { }, { name: "successful true decodes content", - condition: SuccessfulCondition(metav1.ConditionTrue, desire.ReasonSynced), + condition: desiretest.SuccessfulCondition(metav1.ConditionTrue, desire.ReasonSynced), content: configMapManifest(1), wantContent: true, }, { // ConditionTrue/ReasonNotFound is a shape the applier never reports; with no content it reads as not-synced-yet. name: "successful true with empty content is not synced yet", - condition: SuccessfulCondition(metav1.ConditionTrue, desire.ReasonNotFound), + condition: desiretest.SuccessfulCondition(metav1.ConditionTrue, desire.ReasonNotFound), wantNotSynced: true, }, { name: "false with notfound reason is confirmed absent", - condition: SuccessfulCondition(metav1.ConditionFalse, desire.ReasonNotFound), + condition: desiretest.SuccessfulCondition(metav1.ConditionFalse, desire.ReasonNotFound), wantNotFound: true, }, { name: "false with other reason decodes the retained mirror when present", - condition: SuccessfulCondition(metav1.ConditionFalse, desire.ReasonKubeAPIError), + condition: desiretest.SuccessfulCondition(metav1.ConditionFalse, desire.ReasonKubeAPIError), content: configMapManifest(1), wantContent: true, }, { name: "false with other reason and no retained mirror is not synced yet", - condition: SuccessfulCondition(metav1.ConditionFalse, desire.ReasonKubeAPIError), + condition: desiretest.SuccessfulCondition(metav1.ConditionFalse, desire.ReasonKubeAPIError), wantNotSynced: true, }, } diff --git a/internal/executor/resource_executor_test.go b/internal/executor/resource_executor_test.go index 3b3358f4..c531b44a 100644 --- a/internal/executor/resource_executor_test.go +++ b/internal/executor/resource_executor_test.go @@ -7,6 +7,7 @@ import ( "github.com/openshift-hyperfleet/hyperfleet-adapter/internal/configloader" "github.com/openshift-hyperfleet/hyperfleet-adapter/internal/desireclient" + "github.com/openshift-hyperfleet/hyperfleet-adapter/internal/desireclient/desiretest" "github.com/openshift-hyperfleet/hyperfleet-adapter/internal/k8sclient" "github.com/openshift-hyperfleet/hyperfleet-adapter/internal/manifest" "github.com/openshift-hyperfleet/hyperfleet-adapter/internal/transportclient" @@ -2463,7 +2464,7 @@ const ( desireOwner = "hyperfleet-adapter" ) -var testDesireID = desireclient.TestIdentity{ +var testDesireID = desiretest.TestIdentity{ ManagementCluster: "cluster-1", Resource: "configmaps", Namespace: "default", @@ -2542,7 +2543,7 @@ func TestResourceExecutor_DesireTransport_SlowApplier(t *testing.T) { re := newDesireExecutor(store) resource := newDesireResourceWithLifecycle("deleted_time != null", "Background") - desireclient.PutUnsyncedReadDesire(t, ctx, store, testDesireID.Read(), desireOwner) + desiretest.PutUnsyncedReadDesire(t, ctx, store, testDesireID.Read(), desireOwner) // ---- Event 1: apply (deleted_time absent → delete.when false) ---- execCtx := NewExecutionContext(ctx, nil, nil) @@ -2556,7 +2557,7 @@ func TestResourceExecutor_DesireTransport_SlowApplier(t *testing.T) { _, err = store.GetReadDesire(ctx, testDesireID.Read()) require.NoError(t, err, "ReadDesire must exist after apply") - desireclient.MarkReadDesireSynced(t, ctx, store, testDesireID.Read(), configMapContent()) + desiretest.MarkReadDesireSynced(t, ctx, store, testDesireID.Read(), configMapContent()) // ---- Event 2: delete (applier hasn't confirmed yet) ---- execCtx = NewExecutionContext(ctx, nil, nil) @@ -2575,8 +2576,8 @@ func TestResourceExecutor_DesireTransport_SlowApplier(t *testing.T) { _, err = store.GetReadDesire(ctx, testDesireID.Read()) assert.NoError(t, err, "ReadDesire must still exist") - desireclient.MarkDeleteDesireConfirmed(t, ctx, store, testDesireID.Delete()) - desireclient.MarkReadDesireNotFound(t, ctx, store, testDesireID.Read()) + desiretest.MarkDeleteDesireConfirmed(t, ctx, store, testDesireID.Delete()) + desiretest.MarkReadDesireNotFound(t, ctx, store, testDesireID.Read()) // ---- Event 3: delete (applier confirmed) ---- execCtx = NewExecutionContext(ctx, nil, nil) @@ -2608,19 +2609,23 @@ func (s *instantApplierStore) CreateDeleteDesire( if err != nil { return created, err } - s.UpdateDeleteDesireStatus(ctx, dd.Identity, desire.Status{ + if _, err = s.UpdateDeleteDesireStatus(ctx, dd.Identity, desire.Status{ Conditions: []metav1.Condition{{ Type: desire.TypeSuccessful, Status: metav1.ConditionTrue, Reason: desire.ReasonDeleted, }}, - }, created.Version) + }, created.Version); err != nil { + return created, err + } readID := dd.Identity readID.Type = desire.TypeRead - s.UpdateReadDesireStatus(ctx, readID, desire.ReadStatus{ + if _, err = s.UpdateReadDesireStatus(ctx, readID, desire.ReadStatus{ Status: desire.Status{Conditions: []metav1.Condition{{ Type: desire.TypeSuccessful, Status: metav1.ConditionFalse, Reason: desire.ReasonNotFound, }}}, - }) + }); err != nil { + return created, err + } return created, nil } @@ -2643,7 +2648,7 @@ func TestResourceExecutor_DesireTransport_FastApplier(t *testing.T) { }) resource := newDesireResourceWithLifecycle("deleted_time != null", "Background") - desireclient.PutUnsyncedReadDesire(t, ctx, inner, testDesireID.Read(), desireOwner) + desiretest.PutUnsyncedReadDesire(t, ctx, inner, testDesireID.Read(), desireOwner) // ---- Event 1: apply ---- execCtx := NewExecutionContext(ctx, nil, nil) @@ -2652,7 +2657,7 @@ func TestResourceExecutor_DesireTransport_FastApplier(t *testing.T) { require.Len(t, results, 1) assert.Equal(t, StatusSuccess, results[0].Status) - desireclient.MarkReadDesireSynced(t, ctx, inner, testDesireID.Read(), configMapContent()) + desiretest.MarkReadDesireSynced(t, ctx, inner, testDesireID.Read(), configMapContent()) // ---- Event 2: delete (applier confirms instantly via wrapper) ---- execCtx = NewExecutionContext(ctx, nil, nil) @@ -2681,7 +2686,7 @@ func TestResourceExecutor_DesireTransport_TransientNotFound(t *testing.T) { re := newDesireExecutor(store) resource := newDesireResourceWithLifecycle("deleted_time != null", "Background") - desireclient.PutUnsyncedReadDesire(t, ctx, store, testDesireID.Read(), desireOwner) + desiretest.PutUnsyncedReadDesire(t, ctx, store, testDesireID.Read(), desireOwner) // ---- Event 1: apply ---- execCtx := NewExecutionContext(ctx, nil, nil) @@ -2734,7 +2739,7 @@ func TestResourceExecutor_DesireTransport_FullLifecycleFromEmptyStore(t *testing require.NoError(t, err, "ReadDesire must be auto-created by ensureReadDesire") // Simulate applier syncing the read mirror. - desireclient.MarkReadDesireSynced(t, ctx, store, testDesireID.Read(), configMapContent()) + desiretest.MarkReadDesireSynced(t, ctx, store, testDesireID.Read(), configMapContent()) // ---- Event 2: apply after sync (post-apply discovery finds the resource) ---- execCtx = NewExecutionContext(ctx, nil, nil) @@ -2758,8 +2763,8 @@ func TestResourceExecutor_DesireTransport_FullLifecycleFromEmptyStore(t *testing assert.NoError(t, err, "DeleteDesire must exist (pending)") // Simulate applier confirming deletion. - desireclient.MarkDeleteDesireConfirmed(t, ctx, store, testDesireID.Delete()) - desireclient.MarkReadDesireNotFound(t, ctx, store, testDesireID.Read()) + desiretest.MarkDeleteDesireConfirmed(t, ctx, store, testDesireID.Delete()) + desiretest.MarkReadDesireNotFound(t, ctx, store, testDesireID.Read()) // ---- Event 4: delete (applier confirmed → cleanup removes all desires) ---- execCtx = NewExecutionContext(ctx, nil, nil) From 00ca110ab49c12e4b6e996036f287d4a66b14893 Mon Sep 17 00:00:00 2001 From: Michal Vavrinec Date: Tue, 22 Sep 2026 16:53:35 +0200 Subject: [PATCH 6/8] HYPERFLEET-1439 - fix: treat ErrNotSyncedYet as non-fatal during delete discovery Allow delete discovery to proceed through the cleanup path when the desire has not synced yet, instead of failing the execution immediately. - Include ErrNotSyncedYet in the isNotFound classification so cleanup is attempted - Handle ErrDeletionPending in post-delete tryCleanupDesires path - Record deletion error metric only for unexpected errors, not transient states --- internal/executor/resource_executor.go | 23 +++++++++++++---------- 1 file changed, 13 insertions(+), 10 deletions(-) diff --git a/internal/executor/resource_executor.go b/internal/executor/resource_executor.go index 261d750e..4cf186e5 100644 --- a/internal/executor/resource_executor.go +++ b/internal/executor/resource_executor.go @@ -708,17 +708,13 @@ func (re *ResourceExecutor) executeResourceDelete( // Step 1: Discover the existing resource discovered, discoverErr := re.discoverResource(ctx, resource, execCtx, transportClient, transportTarget) - isNotFound := discoverErr != nil && apierrors.IsNotFound(discoverErr) + isNotFound := discoverErr != nil && + (apierrors.IsNotFound(discoverErr) || errors.Is(discoverErr, desireclient.ErrNotSyncedYet)) if discoverErr != nil && !isNotFound { - if errors.Is(discoverErr, desireclient.ErrNotSyncedYet) { - slog.WarnContext(ctx, "resource not synced yet, cannot discover for deletion", - "resource", resource.Name, "error", discoverErr) - } else { - re.metrics.RecordDeletion(resourceType, metrics.DeletionStatusError) - } result.Status = StatusFailed result.Error = discoverErr re.recordResourceError(execCtx, resource, discoverErr) + re.metrics.RecordDeletion(resourceType, metrics.DeletionStatusError) re.metrics.ObserveDeletionDuration(resourceType, time.Since(startTime)) return result, NewExecutorError( PhaseResources, resource.Name, "failed to discover resource for deletion", discoverErr) @@ -730,6 +726,8 @@ func (re *ResourceExecutor) executeResourceDelete( // !resources.?X.hasValue() evaluates to true in this reconciliation. execCtx.Resources[resource.Name] = nil + // Cleanup when discoverErr is ErrNotSyncedYet is safe because ErrDeletionPending is returned when + // the delete desire is still pending. if err := re.tryCleanupDesires(ctx, resource, execCtx, transportClient, transportTarget, gvk); err != nil { if errors.Is(err, desireclient.ErrDeletionPending) { slog.WarnContext(ctx, "resource desire cleanup: deletion pending", @@ -806,12 +804,17 @@ func (re *ResourceExecutor) executeResourceDelete( execCtx.Resources[resource.Name] = nil slog.DebugContext(ctx, "resource confirmed deleted (post-delete discovery: not found)", "resource", resource.Name) if err := re.tryCleanupDesires(ctx, resource, execCtx, transportClient, transportTarget, gvk); err != nil { - slog.ErrorContext(ctx, "resource desire cleanup failed after delete", - "resource", resource.Name, "error", err) + if errors.Is(err, desireclient.ErrDeletionPending) { + slog.WarnContext(ctx, "resource desire cleanup: deletion pending", + "resource", resource.Name, "error", err) + } else { + slog.ErrorContext(ctx, "resource desire cleanup failed after delete", + "resource", resource.Name, "error", err) + re.metrics.RecordDeletion(resourceType, metrics.DeletionStatusError) + } result.Status = StatusFailed result.Error = err re.recordResourceError(execCtx, resource, err) - re.metrics.RecordDeletion(resourceType, metrics.DeletionStatusError) re.metrics.ObserveDeletionDuration(resourceType, time.Since(startTime)) return result, NewExecutorError(PhaseResources, resource.Name, "desire cleanup failed", err) } From 4e423a0595960e5744a64fff46b6d1fd8d6677e8 Mon Sep 17 00:00:00 2001 From: Michal Vavrinec Date: Wed, 23 Sep 2026 10:14:28 +0200 Subject: [PATCH 7/8] HYPERFLEET-1439 - fix: treat ErrDeletionPending as non-fatal during post-delete cleanup Previously, ErrDeletionPending at Step 6 caused StatusFailed. Now Step 6 restores the pre-delete discovered state in context so dependents wait, and continues to success. - Make ErrDeletionPending non-fatal at Step 6, restore discovered state for dependents - Add Step 2 and Step 6 ErrDeletionPending tests with real desire store - Move InstantApplierStore and PendingDeleteApplierStore to desiretest package - Widen newDesireExecutor to accept desire.SpecStore for better reusability --- .../desireclient/desiretest/desiretest.go | 62 ++++++++ internal/executor/resource_executor.go | 22 +-- internal/executor/resource_executor_test.go | 150 ++++++++++++------ 3 files changed, 176 insertions(+), 58 deletions(-) diff --git a/internal/desireclient/desiretest/desiretest.go b/internal/desireclient/desiretest/desiretest.go index 242f723f..0aff86c5 100644 --- a/internal/desireclient/desiretest/desiretest.go +++ b/internal/desireclient/desiretest/desiretest.go @@ -199,3 +199,65 @@ func (ti TestIdentity) WithNamespace(namespace string) TestIdentity { ti.Namespace = namespace return ti } + +// InstantApplierStore wraps memory.Store. When a DeleteDesire is created it +// immediately marks it Deleted and the paired ReadDesire as NotFound, +// simulating an applier that confirms before post-delete discovery runs. +type InstantApplierStore struct { + *memory.Store +} + +func (s *InstantApplierStore) CreateDeleteDesire( + ctx context.Context, dd desire.DeleteDesire, +) (desire.DeleteDesire, error) { + created, err := s.Store.CreateDeleteDesire(ctx, dd) + if err != nil { + return created, err + } + if _, err = s.UpdateDeleteDesireStatus(ctx, dd.Identity, desire.Status{ + Conditions: []metav1.Condition{{ + Type: desire.TypeSuccessful, Status: metav1.ConditionTrue, Reason: desire.ReasonDeleted, + }}, + }, created.Version); err != nil { + return created, err + } + + readID := dd.Identity + readID.Type = desire.TypeRead + if _, err = s.UpdateReadDesireStatus(ctx, readID, desire.ReadStatus{ + Status: desire.Status{Conditions: []metav1.Condition{{ + Type: desire.TypeSuccessful, Status: metav1.ConditionFalse, Reason: desire.ReasonNotFound, + }}}, + }); err != nil { + return created, err + } + return created, nil +} + +// PendingDeleteApplierStore wraps memory.Store. When a DeleteDesire is +// created it marks the paired ReadDesire as NotFound (applier saw the +// resource gone) but does NOT confirm the delete desire, simulating +// an applier that is slow to ack the deletion. +type PendingDeleteApplierStore struct { + *memory.Store +} + +func (s *PendingDeleteApplierStore) CreateDeleteDesire( + ctx context.Context, dd desire.DeleteDesire, +) (desire.DeleteDesire, error) { + created, err := s.Store.CreateDeleteDesire(ctx, dd) + if err != nil { + return created, err + } + + readID := dd.Identity + readID.Type = desire.TypeRead + if _, err = s.UpdateReadDesireStatus(ctx, readID, desire.ReadStatus{ + Status: desire.Status{Conditions: []metav1.Condition{{ + Type: desire.TypeSuccessful, Status: metav1.ConditionFalse, Reason: desire.ReasonNotFound, + }}}, + }); err != nil { + return created, err + } + return created, nil +} diff --git a/internal/executor/resource_executor.go b/internal/executor/resource_executor.go index 4cf186e5..068da15f 100644 --- a/internal/executor/resource_executor.go +++ b/internal/executor/resource_executor.go @@ -802,21 +802,23 @@ func (re *ResourceExecutor) executeResourceDelete( case postDeleteDiscovered == nil || postIsNotFound: // Resource is confirmed gone: dependent resources can proceed in this reconciliation. execCtx.Resources[resource.Name] = nil - slog.DebugContext(ctx, "resource confirmed deleted (post-delete discovery: not found)", "resource", resource.Name) + slog.DebugContext(ctx, "post-delete discovery: resource not found)", "resource", resource.Name) if err := re.tryCleanupDesires(ctx, resource, execCtx, transportClient, transportTarget, gvk); err != nil { - if errors.Is(err, desireclient.ErrDeletionPending) { - slog.WarnContext(ctx, "resource desire cleanup: deletion pending", - "resource", resource.Name, "error", err) - } else { + if !errors.Is(err, desireclient.ErrDeletionPending) { slog.ErrorContext(ctx, "resource desire cleanup failed after delete", "resource", resource.Name, "error", err) + result.Status = StatusFailed + result.Error = err + re.recordResourceError(execCtx, resource, err) re.metrics.RecordDeletion(resourceType, metrics.DeletionStatusError) + re.metrics.ObserveDeletionDuration(resourceType, time.Since(startTime)) + return result, NewExecutorError(PhaseResources, resource.Name, "desire cleanup failed", err) } - result.Status = StatusFailed - result.Error = err - re.recordResourceError(execCtx, resource, err) - re.metrics.ObserveDeletionDuration(resourceType, time.Since(startTime)) - return result, NewExecutorError(PhaseResources, resource.Name, "desire cleanup failed", err) + + // Deletion pending: leave the last known discovered state in context + execCtx.Resources[resource.Name] = discovered + slog.WarnContext(ctx, "resource desire cleanup: deletion pending", + "resource", resource.Name, "error", err) } default: // Resource still present (finalizers or async deletion): dependents wait for next reconciliation. diff --git a/internal/executor/resource_executor_test.go b/internal/executor/resource_executor_test.go index c531b44a..74f11bab 100644 --- a/internal/executor/resource_executor_test.go +++ b/internal/executor/resource_executor_test.go @@ -16,7 +16,6 @@ import ( "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" "k8s.io/apimachinery/pkg/apis/meta/v1/unstructured" "k8s.io/apimachinery/pkg/runtime/schema" ) @@ -2484,7 +2483,7 @@ func configMapContent() []byte { }`) } -func newDesireExecutor(store *memory.Store) *ResourceExecutor { +func newDesireExecutor(store desire.SpecStore) *ResourceExecutor { c := desireclient.NewClient(store, desireOwner) return newResourceExecutor(&ExecutorConfig{ Config: &configloader.Config{ @@ -2595,57 +2594,13 @@ func TestResourceExecutor_DesireTransport_SlowApplier(t *testing.T) { assert.ErrorIs(t, err, desire.ErrNotFound, "ReadDesire must be removed by cleanup") } -// instantApplierStore wraps memory.Store. When a DeleteDesire is created it -// immediately marks it Deleted and the paired ReadDesire as NotFound, -// simulating an applier that confirms before post-delete discovery runs. -type instantApplierStore struct { - *memory.Store -} - -func (s *instantApplierStore) CreateDeleteDesire( - ctx context.Context, dd desire.DeleteDesire, -) (desire.DeleteDesire, error) { - created, err := s.Store.CreateDeleteDesire(ctx, dd) - if err != nil { - return created, err - } - if _, err = s.UpdateDeleteDesireStatus(ctx, dd.Identity, desire.Status{ - Conditions: []metav1.Condition{{ - Type: desire.TypeSuccessful, Status: metav1.ConditionTrue, Reason: desire.ReasonDeleted, - }}, - }, created.Version); err != nil { - return created, err - } - - readID := dd.Identity - readID.Type = desire.TypeRead - if _, err = s.UpdateReadDesireStatus(ctx, readID, desire.ReadStatus{ - Status: desire.Status{Conditions: []metav1.Condition{{ - Type: desire.TypeSuccessful, Status: metav1.ConditionFalse, Reason: desire.ReasonNotFound, - }}}, - }); err != nil { - return created, err - } - return created, nil -} - // Test 2: Fast applier — applier confirms between DeleteResource and // post-delete discovery within the same event. func TestResourceExecutor_DesireTransport_FastApplier(t *testing.T) { ctx := context.Background() inner := memory.New() - store := &instantApplierStore{Store: inner} - c := desireclient.NewClient(store, desireOwner) - re := newResourceExecutor(&ExecutorConfig{ - Config: &configloader.Config{ - Transports: map[string]configloader.TransportDefinition{ - desireTransportName: {Type: configloader.TransportTypeRemote}, - }, - }, - TransportRegistry: transportclient.Registry{ - desireTransportName: c, - }, - }) + store := &desiretest.InstantApplierStore{Store: inner} + re := newDesireExecutor(store) resource := newDesireResourceWithLifecycle("deleted_time != null", "Background") desiretest.PutUnsyncedReadDesire(t, ctx, inner, testDesireID.Read(), desireOwner) @@ -2792,3 +2747,102 @@ func TestResourceExecutor_DesireTransport_FullLifecycleFromEmptyStore(t *testing _, err = store.GetReadDesire(ctx, testDesireID.Read()) require.NoError(t, err, "ReadDesire must be re-created after cleanup") } + +// Test 5: Step 2 ErrDeletionPending — delete desire exists but applier +// hasn't confirmed yet. The read desire shows NotFound (applier's +// informer observed the resource gone) but the delete desire is still +// pending. Cleanup must refuse with ErrDeletionPending, keeping both +// desires alive so the next reconciliation retries. +func TestResourceExecutor_DesireTransport_Step2_DeletionPending_DeleteNotConfirmed(t *testing.T) { + ctx := context.Background() + store := memory.New() + re := newDesireExecutor(store) + resource := newDesireResourceWithLifecycle("deleted_time != null", "Background") + + desiretest.PutUnsyncedReadDesire(t, ctx, store, testDesireID.Read(), desireOwner) + + // ---- Event 1: apply ---- + execCtx := NewExecutionContext(ctx, nil, nil) + results, err := re.ExecuteAll(ctx, []configloader.Resource{resource}, execCtx) + require.NoError(t, err) + require.Len(t, results, 1) + assert.Equal(t, StatusSuccess, results[0].Status) + + desiretest.MarkReadDesireSynced(t, ctx, store, testDesireID.Read(), configMapContent()) + + // ---- Event 2: delete (creates delete desire, applier hasn't confirmed) ---- + execCtx = NewExecutionContext(ctx, nil, nil) + execCtx.Params["deleted_time"] = testDeletedTime + results, err = re.ExecuteAll(ctx, []configloader.Resource{resource}, execCtx) + require.NoError(t, err) + require.Len(t, results, 1) + assert.Equal(t, StatusSuccess, results[0].Status) + + _, err = store.GetDeleteDesire(ctx, testDesireID.Delete()) + require.NoError(t, err, "DeleteDesire must exist (pending)") + + // Applier's read informer saw the resource gone but the delete + // desire hasn't been confirmed yet. + desiretest.MarkReadDesireNotFound(t, ctx, store, testDesireID.Read()) + + // ---- Event 3: delete (discovery → NotFound → Step 2 → cleanup → ErrDeletionPending) ---- + execCtx = NewExecutionContext(ctx, nil, nil) + execCtx.Params["deleted_time"] = testDeletedTime + results, err = re.ExecuteAll(ctx, []configloader.Resource{resource}, execCtx) + require.Error(t, err) + assert.ErrorIs(t, results[0].Error, desireclient.ErrDeletionPending, + "cleanup must return ErrDeletionPending when delete desire is unconfirmed") + assert.Equal(t, StatusFailed, results[0].Status) + + _, err = store.GetDeleteDesire(ctx, testDesireID.Delete()) + assert.NoError(t, err, "DeleteDesire must survive — not yet confirmed") + + _, err = store.GetReadDesire(ctx, testDesireID.Read()) + assert.NoError(t, err, "ReadDesire must survive") +} + +// Test 6: Step 6 ErrDeletionPending — resource existed, was deleted, +// post-delete discovery confirms gone, but the delete desire is still +// pending. Cleanup returns ErrDeletionPending which is non-fatal at +// Step 6: the executor restores the pre-delete discovered state in +// context (dependents wait) and still reports success. +func TestResourceExecutor_DesireTransport_Step6_DeletionPending_NonFatal(t *testing.T) { + ctx := context.Background() + inner := memory.New() + store := &desiretest.PendingDeleteApplierStore{Store: inner} + re := newDesireExecutor(store) + resource := newDesireResourceWithLifecycle("deleted_time != null", "Background") + + desiretest.PutUnsyncedReadDesire(t, ctx, inner, testDesireID.Read(), desireOwner) + + // ---- Event 1: apply ---- + execCtx := NewExecutionContext(ctx, nil, nil) + results, err := re.ExecuteAll(ctx, []configloader.Resource{resource}, execCtx) + require.NoError(t, err) + require.Len(t, results, 1) + assert.Equal(t, StatusSuccess, results[0].Status) + + desiretest.MarkReadDesireSynced(t, ctx, inner, testDesireID.Read(), configMapContent()) + + // ---- Event 2: delete ---- + // Discovery finds the resource (synced). DeleteResource creates a + // delete desire. The pendingDeleteApplierStore wrapper marks the + // read desire as NotFound (applier saw it gone) but does NOT confirm + // the delete desire. Post-delete discovery → NotFound → Step 6 → + // tryCleanupDesires → delete desire not confirmed → ErrDeletionPending. + // Step 6 treats ErrDeletionPending as non-fatal. + execCtx = NewExecutionContext(ctx, nil, nil) + execCtx.Params["deleted_time"] = testDeletedTime + results, err = re.ExecuteAll(ctx, []configloader.Resource{resource}, execCtx) + require.NoError(t, err, "ErrDeletionPending at Step 6 must be non-fatal") + require.Len(t, results, 1) + assert.Equal(t, StatusSuccess, results[0].Status) + assert.NotNil(t, execCtx.Resources["test-resource"], + "pre-delete discovered state must be restored so dependents wait") + + _, err = inner.GetDeleteDesire(ctx, testDesireID.Delete()) + assert.NoError(t, err, "DeleteDesire must still exist (pending, not confirmed)") + + _, err = inner.GetApplyDesire(ctx, testDesireID.Apply()) + assert.ErrorIs(t, err, desire.ErrNotFound, "ApplyDesire must be removed by DeleteResource") +} From 8b311629d3f690af07f91c9a835e1309c0065f8c Mon Sep 17 00:00:00 2001 From: Michal Vavrinec Date: Wed, 23 Sep 2026 11:47:09 +0200 Subject: [PATCH 8/8] HYPERFLEET-1439 - docs: clarify error recording in post-delete cleanup Adds a comment explaining why the result status and error are recorded after a post-delete action failure. - Document that Health=False reporting depends on these fields being set - Note the constraint: Finalized=True cannot be prevented without exposing a discovered resource in the execution context --- internal/executor/resource_executor.go | 3 +++ 1 file changed, 3 insertions(+) diff --git a/internal/executor/resource_executor.go b/internal/executor/resource_executor.go index 068da15f..15469cad 100644 --- a/internal/executor/resource_executor.go +++ b/internal/executor/resource_executor.go @@ -737,6 +737,9 @@ func (re *ResourceExecutor) executeResourceDelete( "resource", resource.Name, "error", err) re.metrics.RecordDeletion(resourceType, metrics.DeletionStatusError) } + // Record the error in the result and execution context so Health=False is reported + // since we cannot prevent Finalized=True without exposing a discovered resource in + // the context result.Status = StatusFailed result.Error = err re.recordResourceError(execCtx, resource, err)