diff --git a/internal/desireclient/cleanup.go b/internal/desireclient/cleanup.go new file mode 100644 index 00000000..a31b9c2b --- /dev/null +++ b/internal/desireclient/cleanup.go @@ -0,0 +1,91 @@ +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, +// 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, + 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): + 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: %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", + 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..89d21567 --- /dev/null +++ b/internal/desireclient/cleanup_test.go @@ -0,0 +1,269 @@ +package desireclient + +import ( + "context" + "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" + 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, + } + + desiretest.PutConfirmedDeleteDesire(t, ctx, store, testID.Delete(), testOwner) + + _, 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, + } + + desiretest.PutDeleteDesire(t, ctx, store, testID.Delete(), testOwner, + 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") + assert.True(t, errors.Is(err, ErrDeletionPending), "must wrap ErrDeletionPending") + + _, 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_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() + c := newTestClient(store) + + deleteID := desire.Identity{ + ManagementCluster: testManagementCluster, Type: desire.TypeDelete, + Resource: testResource, Namespace: testNamespace, Name: testName, + } + + desiretest.PutConfirmedDeleteDesire(t, ctx, store, testID.Delete(), testOwner) + + 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() + + desiretest.PutConfirmedDeleteDesire(t, ctx, inner, testID.Delete(), testOwner) + + 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") +} + +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 { + desire.SpecStore +} + +func (f *failingGetDeleteDesireStore) GetDeleteDesire( + _ context.Context, _ desire.Identity, +) (desire.DeleteDesire, error) { + 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 +} + +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/desireclient_test.go b/internal/desireclient/desireclient_test.go index 80f8e774..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,6 +73,54 @@ func (f *failingCreateDeleteDesireStore) CreateDeleteDesire( return desire.DeleteDesire{}, errors.New("boom: delete desire store unavailable") } +var testID = desiretest.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/desiretest/desiretest.go b/internal/desireclient/desiretest/desiretest.go new file mode 100644 index 00000000..0aff86c5 --- /dev/null +++ b/internal/desireclient/desiretest/desiretest.go @@ -0,0 +1,263 @@ +package desiretest + +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" +) + +// 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, + 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 +} + +// 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/desireclient/discover_test.go b/internal/desireclient/discover_test.go index f5911c68..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, testNamespace, testName, 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, testNamespace, testName, 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,8 +60,10 @@ 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) + desiretest.PutSyncedReadDesire(t, ctx, store, + testID.WithNamespace("labeled").WithName("labeled").Read(), testOwner, labeledManifest) + desiretest.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 +94,7 @@ func TestDiscoverResources_SurfacesRetainedMirrorOnFailedRead(t *testing.T) { store := newMemoryStore() c := newTestClient(store) - putKubeAPIErrorReadDesire(t, ctx, store, testNamespace, testName, 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) @@ -105,7 +108,7 @@ func TestDiscoverResources_SkipsFailedReadWithNoRetainedMirror(t *testing.T) { store := newMemoryStore() c := newTestClient(store) - putKubeAPIErrorReadDesire(t, ctx, store, testNamespace, testName, nil) + desiretest.PutKubeAPIErrorReadDesire(t, ctx, store, testID.Read(), testOwner, nil) list, err := c.DiscoverResources(ctx, testGVK(), &manifest.DiscoveryConfig{}, testTransportContext()) require.NoError(t, err) @@ -117,7 +120,7 @@ func TestDiscoverResources_SkipsNotFoundFalseDesire(t *testing.T) { store := newMemoryStore() c := newTestClient(store) - putNotFoundReadDesire(t, ctx, store, testNamespace, testName) + desiretest.PutNotFoundReadDesire(t, ctx, store, testID.Read(), testOwner) list, err := c.DiscoverResources(ctx, testGVK(), &manifest.DiscoveryConfig{}, testTransportContext()) require.NoError(t, err) @@ -129,8 +132,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)) + 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") @@ -143,7 +146,7 @@ func TestDiscoverResources_FiltersOutOtherResourceType(t *testing.T) { store := newMemoryStore() c := newTestClient(store) - putSyncedReadDesire(t, ctx, store, testNamespace, testName, 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) @@ -182,7 +185,7 @@ func TestDiscoverResources_ScopedToPartition(t *testing.T) { store := newMemoryStore() c := newTestClient(store) - putSyncedReadDesire(t, ctx, store, testNamespace, testName, 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) @@ -197,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, testNamespace, "first", first) - putSyncedReadDesire(t, ctx, store, testNamespace, "second", 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 1804aa81..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" @@ -12,13 +13,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 +29,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 +43,7 @@ func TestGetResource_SyncedReturnsMirroredObject(t *testing.T) { store := newMemoryStore() c := newTestClient(store) - putSyncedReadDesire(t, ctx, store, testNamespace, testName, 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) @@ -62,7 +56,7 @@ func TestGetResource_ConfirmedNotFound(t *testing.T) { store := newMemoryStore() c := newTestClient(store) - putConfirmedAbsentReadDesire(t, ctx, store, testNamespace, testName) + desiretest.PutConfirmedAbsentReadDesire(t, ctx, store, testID.Read(), testOwner) _, err := c.GetResource(ctx, testGVK(), testNamespace, testName, testTransportContext()) require.Error(t, err) @@ -75,7 +69,7 @@ func TestGetResource_InvalidReadDesire(t *testing.T) { store := newMemoryStore() c := newTestClient(store) - putInvalidReadDesire(t, ctx, store, testNamespace, testName) + desiretest.PutInvalidReadDesire(t, ctx, store, testID.Read(), testOwner) _, err := c.GetResource(ctx, testGVK(), testNamespace, testName, testTransportContext()) require.Error(t, err) @@ -90,7 +84,7 @@ func TestGetResource_K8sAPIErrorWithRetainedMirrorReturnsStaleContent(t *testing store := newMemoryStore() c := newTestClient(store) - putKubeAPIErrorReadDesire(t, ctx, store, testNamespace, testName, 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, @@ -104,7 +98,7 @@ func TestGetResource_K8sAPIErrorWithNoMirrorYetIsNotSyncedYet(t *testing.T) { store := newMemoryStore() c := newTestClient(store) - putKubeAPIErrorReadDesire(t, ctx, store, testNamespace, testName, nil) + desiretest.PutKubeAPIErrorReadDesire(t, ctx, store, testID.Read(), testOwner, nil) _, err := c.GetResource(ctx, testGVK(), testNamespace, testName, testTransportContext()) require.Error(t, err) @@ -159,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, }, } @@ -190,7 +184,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 e5768b13..00000000 --- a/internal/desireclient/helpers_test.go +++ /dev/null @@ -1,129 +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, - }) -} - -// successfulCondition builds the single summary condition every desire carries. -func successfulCondition(status metav1.ConditionStatus, reason string) *metav1.Condition { - return &metav1.Condition{Type: desire.TypeSuccessful, Status: status, Reason: reason} -} - -// failingListReadDesiresStore wraps a real SpecStore but forces -// ListReadDesires to fail, simulating a store-level outage during discovery. -type failingListReadDesiresStore struct { - desire.SpecStore -} - -func (f *failingListReadDesiresStore) ListReadDesires( - _ context.Context, _ string, -) ([]desire.ReadDesire, error) { - return nil, errors.New("boom: read desire store unavailable") -} - -// spyDeleteReadDesireStore counts DeleteReadDesire calls so tests can assert -// whether ensureReadDesire actually attempted a recreate. -type spyDeleteReadDesireStore struct { - desire.SpecStore - deleteReadDesireCalls int -} - -func (s *spyDeleteReadDesireStore) DeleteReadDesire( - ctx context.Context, id desire.Identity, owner string, version int64, -) error { - s.deleteReadDesireCalls++ - return s.SpecStore.DeleteReadDesire(ctx, id, owner, version) -} - -// staleApplyVersionStore wraps a real SpecStore but returns a stale version -// on GetApplyDesire to simulate the case where an external client has -// concurrently updated the apply desire while this one is computing. -type staleApplyVersionStore struct { - desire.SpecStore -} - -func (s *staleApplyVersionStore) GetApplyDesire(ctx context.Context, id desire.Identity) (desire.ApplyDesire, error) { - ad, err := s.SpecStore.GetApplyDesire(ctx, id) - if err == nil { - ad.Version++ - } - return ad, err -} diff --git a/internal/desireclient/types.go b/internal/desireclient/types.go 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 5186849a..15469cad 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 @@ -291,6 +296,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 +332,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 +360,7 @@ func (re *ResourceExecutor) discoverResource( labelSelector := manifest.BuildLabelSelector(renderedLabels) discoveryConfig := &manifest.DiscoveryConfig{ - Namespace: namespace, + Namespace: dt.Namespace, LabelSelector: labelSelector, } @@ -361,6 +381,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( @@ -526,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 } @@ -669,7 +708,8 @@ 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 { result.Status = StatusFailed result.Error = discoverErr @@ -685,8 +725,30 @@ 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" + + // 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", + "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) + } + // 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) + 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 @@ -743,7 +805,24 @@ 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.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) + } + + // 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. execCtx.Resources[resource.Name] = postDeleteDiscovered diff --git a/internal/executor/resource_executor_test.go b/internal/executor/resource_executor_test.go index 11d19bac..74f11bab 100644 --- a/internal/executor/resource_executor_test.go +++ b/internal/executor/resource_executor_test.go @@ -7,9 +7,12 @@ 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" + "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" @@ -2284,3 +2287,562 @@ 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") +} + +// ---- Desire transport lifecycle tests ---- + +const ( + desireTransportName = "desire-primary" + desireOwner = "hyperfleet-adapter" +) + +var testDesireID = desiretest.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 desire.SpecStore) *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") + + desiretest.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") + + desiretest.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") + + desiretest.MarkDeleteDesireConfirmed(t, ctx, store, testDesireID.Delete()) + desiretest.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") +} + +// 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 := &desiretest.InstantApplierStore{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 (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") + + 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) + + // 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. + desiretest.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. + 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) + 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") +} + +// 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") +} 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 +}