-
Notifications
You must be signed in to change notification settings - Fork 23
HYPERFLEET-1439 - feat: cleanup desires after deletion #303
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Changes from all commits
8676a9d
a473799
2f67f45
ea8e558
63e5ab4
00ca110
4e423a0
8b31162
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -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): | ||
|
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. So coderabbit picked up on this from a concurrency angle, a reapply racing cleanup. But it also can fire sequentially too, so it is definitely a race we want to patch. Help paint that picture ill walk through two events for the same cluster :
|
||
| 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 { | ||
|
coderabbitai[bot] marked this conversation as resolved.
|
||
| 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 { | ||
|
Comment on lines
+74
to
+82
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. 🗄️ Data Integrity & Integration | 🟠 Major | 🏗️ Heavy lift 🔎 Supported by static analysis🏁 Script executed: sed -n '1,120p' internal/desireclient/cleanup.go
rg -n -C 3 'CreateReadDesire|ensureReadDesire|DeleteDeleteDesire|CleanupAfterDeletion|CreateDeleteDesire' internal cmd
rg -n -C 3 'Subscribe|handler|goroutine|parallel|concurr' cmd internalRepository: openshift-hyperfleet/hyperfleet-adapter Length of output: 50397 🏁 Script executed: #!/bin/bash
set -u
printf '%s\n' '--- cleanup and ensure ---'
sed -n '1,190p' internal/desireclient/cleanup.go
sed -n '120,185p' internal/desireclient/apply.go
printf '%s\n' '--- executor entry and relevant lifecycle ---'
sed -n '1,180p' internal/executor/handler.go
sed -n '330,405p' internal/executor/resource_executor.go
sed -n '500,565p' cmd/adapter/main.go
printf '%s\n' '--- broker subscriber binding ---'
rg -n -C 5 'type Subscriber|func .*Subscribe|parallel|goroutine|worker|handler' "$(go env GOPATH 2>/dev/null)/pkg/mod/github.com/openshift-hyperfleet" 2>/dev/null || true
rg -n -C 4 'Subscriber|Subscribe' go.mod go.sum internal cmd
printf '%s\n' '--- desire store binding and CAS behavior ---'
rg -n -C 5 'type SpecStore|DeleteReadDesire|GetReadDesire|CreateReadDesire' . --glob '*.go' --glob '!**/*_test.go'Repository: openshift-hyperfleet/hyperfleet-adapter Length of output: 19826 🏁 Script executed: #!/bin/bash
set -u
printf '%s\n' '--- subscriber and adapter binding ---'
rg -n -C 6 'type Subscriber|func .*Subscribe|Subscribe\(' . --glob '*.go' --glob '!**/*_test.go' | head -240
printf '%s\n' '--- executor invocation ---'
rg -n -C 6 'CreateHandler|func \(.*\) Execute|Execute\(' internal/executor cmd --glob '*.go' | head -260
printf '%s\n' '--- desire store calls ---'
rg -n -C 8 'type SpecStore|DeleteReadDesire|GetReadDesire|CreateReadDesire' . --glob '*.go' --glob '!**/*_test.go' | head -360
printf '%s\n' '--- ensureReadDesire ---'
sed -n '125,180p' internal/desireclient/apply.goRepository: openshift-hyperfleet/hyperfleet-adapter Length of output: 27524 🏁 Script executed: #!/bin/bash
set -u
rg -n -C 8 'type Subscriber|func .*Subscribe|parallel|goroutine|worker|handler' . --glob '*.go' | head -300Repository: openshift-hyperfleet/hyperfleet-broker Length of output: 20151 🏁 Script executed: #!/bin/bash
set -u
printf '%s\n' '--- ApplyResource ordering and delete-intent handling ---'
sed -n '1,125p' internal/desireclient/apply.go
printf '%s\n' '--- store interface and local implementations ---'
rg -n -C 10 'type SpecStore|func .*DeleteReadDesire|func .*GetReadDesire|func .*CreateReadDesire' . --glob '*.go' --glob '!**/*_test.go' | head -420Repository: openshift-hyperfleet/hyperfleet-adapter Length of output: 5017 Make desire cleanup atomic to avoid deleting a replacement read desire. 🤖 Prompt for AI Agents |
||
| 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 | ||
| } | ||
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -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") | ||
| } |
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
📐 Maintainability & Code Quality | 🟡 Minor | ⚡ Quick win
Wrap each returned error with cleanup context.
These bare returns lose the failed cleanup stage. Add context for transport resolution and identity construction.
Proposed fix
tc, err := resolveTransportContext(target) if err != nil { - return err + return fmt.Errorf("desireclient: cleanup: resolve transport context: %w", err) } deleteID, err := buildIdentity(tc, desire.TypeDelete, gvk, namespace, name) if err != nil { - return err + return fmt.Errorf("desireclient: cleanup: build delete desire identity: %w", err) } ... readID, err := buildIdentity(tc, desire.TypeRead, gvk, namespace, name) if err != nil { - return err + return fmt.Errorf("desireclient: cleanup: build read desire identity: %w", err) }As per path instructions, “Wrap errors per Error Model Standard — no bare return err.”
Also applies to: 31-31, 55-55
🤖 Prompt for AI Agents
Source: Path instructions