Skip to content

Commit 2fbce2e

Browse files
authored
fix(controller): skip reconcile for tasks that are being deleted (#417)
1 parent c5c1ac5 commit 2fbce2e

2 files changed

Lines changed: 64 additions & 0 deletions

File tree

‎internal/controller/worker.go‎

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -126,6 +126,12 @@ func (w *Worker) processEvent(ctx context.Context, ev store.TaskEvent) error {
126126
}
127127
return fmt.Errorf("fetching task %s/%s: %w", ev.Atespace, ev.Name, err)
128128
}
129+
// A pending delete event owns this task now; reconciling would resume an actor
130+
// that is about to be torn down and overwrite the Terminating phase.
131+
if task.GetStatus().GetPhase() == v1alpha1.PhaseTerminating {
132+
slog.Info("task is terminating, skipping reconcile", "atespace", ev.Atespace, "name", ev.Name)
133+
return nil
134+
}
129135

130136
// Resolve every bound workspace. A missing one is skipped so the task still
131137
// runs; the runner creates an empty directory at its path.

‎internal/controller/worker_test.go‎

Lines changed: 58 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -173,3 +173,61 @@ func TestWorkerDeletion(t *testing.T) {
173173
t.Errorf("expected template deleted, got %v", mockSrv.deletedTemplates)
174174
}
175175
}
176+
177+
func TestWorkerSkipsReconcileOfTerminatingTask(t *testing.T) {
178+
ctx, cancel := context.WithCancel(context.Background())
179+
defer cancel()
180+
181+
lis, err := net.Listen("tcp", "127.0.0.1:0")
182+
if err != nil {
183+
t.Fatalf("failed to listen: %v", err)
184+
}
185+
defer lis.Close()
186+
187+
mockSrv := &mockControlServer{}
188+
grpcServer := grpc.NewServer()
189+
ateapipb.RegisterControlServer(grpcServer, mockSrv)
190+
go grpcServer.Serve(lis)
191+
defer grpcServer.Stop()
192+
193+
subClient, err := substrate.NewClient(lis.Addr().String(), grpc.WithTransportCredentials(insecure.NewCredentials()))
194+
if err != nil {
195+
t.Fatalf("failed to create substrate client: %v", err)
196+
}
197+
defer subClient.Close()
198+
199+
reconciler := controller.NewTaskReconciler(subClient, "default-template", "ax-system")
200+
reconciler.SecretResolver = noSecrets
201+
reconciler.WorkspaceReadyTimeout = 200 * time.Millisecond
202+
203+
// Queue a reconcile and then a delete before the worker starts, as when a
204+
// task is deleted while the controller is still busy with other events.
205+
memStore := memory.NewStore()
206+
task := &v1alpha1.Task{
207+
Metadata: &v1alpha1.ObjectMeta{Name: "doomed", Atespace: "default"},
208+
Spec: &v1alpha1.TaskSpec{Image: "ghcr.io/test/img"},
209+
}
210+
if err := memStore.SaveTask(ctx, task); err != nil {
211+
t.Fatalf("failed to save task: %v", err)
212+
}
213+
if err := memStore.MarkTaskDeleting(ctx, "default", "doomed"); err != nil {
214+
t.Fatalf("MarkTaskDeleting failed: %v", err)
215+
}
216+
217+
worker := controller.NewWorker(memStore, reconciler, "test-group", "worker-1")
218+
go func() { _ = worker.Run(ctx) }()
219+
220+
deadline := time.Now().Add(3 * time.Second)
221+
for time.Now().Before(deadline) {
222+
if _, err := memStore.GetTask(ctx, "default", "doomed"); err != nil {
223+
break
224+
}
225+
time.Sleep(50 * time.Millisecond)
226+
}
227+
if _, err := memStore.GetTask(ctx, "default", "doomed"); err == nil {
228+
t.Fatalf("expected task record to be removed after cleanup")
229+
}
230+
if len(mockSrv.createdActors) != 0 || len(mockSrv.resumedActors) != 0 {
231+
t.Errorf("terminating task was reconciled: created %v, resumed %v", mockSrv.createdActors, mockSrv.resumedActors)
232+
}
233+
}

0 commit comments

Comments
 (0)