Skip to content

Commit 07451e4

Browse files
committed
[task-controller] Task Controller supports parallel reconciles
1 parent 72a423d commit 07451e4

2 files changed

Lines changed: 39 additions & 24 deletions

File tree

control-operator/cmd/task-manager/main.go

Lines changed: 8 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -57,11 +57,14 @@ func main() {
5757
var metricsAddr string
5858
var enableLeaderElection bool
5959
var probeAddr string
60+
var maxConcurrentReconciles int
6061
flag.StringVar(&metricsAddr, "metrics-bind-address", ":9082", "The address the metric endpoint binds to.")
6162
flag.StringVar(&probeAddr, "health-probe-bind-address", ":9083", "The address the probe endpoint binds to.")
6263
flag.BoolVar(&enableLeaderElection, "leader-elect", false,
6364
"Enable leader election for controller manager. "+
6465
"Enabling this will ensure there is only one active controller manager.")
66+
flag.IntVar(&maxConcurrentReconciles, "max-concurrent-reconciles", 1,
67+
"The maximum number of concurrent Reconciles which can be run for the Task controller.")
6568
opts := zap.Options{
6669
Development: true,
6770
}
@@ -90,10 +93,11 @@ func main() {
9093
}
9194

9295
if err = (&controller.TaskReconciler{
93-
Client: mgr.GetClient(),
94-
Scheme: mgr.GetScheme(),
95-
Recorder: mgr.GetEventRecorderFor("task-controller"),
96-
NodeName: nodeName,
96+
Client: mgr.GetClient(),
97+
Scheme: mgr.GetScheme(),
98+
Recorder: mgr.GetEventRecorderFor("task-controller"),
99+
NodeName: nodeName,
100+
MaxConcurrentReconciles: maxConcurrentReconciles,
97101
}).SetupWithManager(mgr); err != nil {
98102
setupLog.Error(err, "unable to create controller", "controller", "Task")
99103
os.Exit(1)

control-operator/internal/controller/task_controller.go

Lines changed: 31 additions & 20 deletions
Original file line numberDiff line numberDiff line change
@@ -30,6 +30,7 @@ import (
3030
"fmt"
3131
"reflect"
3232
"strings"
33+
"sync"
3334
"time"
3435

3536
v1 "k8s.io/api/core/v1"
@@ -54,12 +55,14 @@ import (
5455
// TaskReconciler reconciles a Task object
5556
type TaskReconciler struct {
5657
client.Client
57-
Scheme *runtime.Scheme
58-
Recorder record.EventRecorder
59-
NodeName string
58+
Scheme *runtime.Scheme
59+
Recorder record.EventRecorder
60+
NodeName string
61+
MaxConcurrentReconciles int
6062
}
6163

62-
var clientsForContainers map[string]*OccClient = make(map[string]*OccClient)
64+
// var clientsForContainers map[string]*OccClient = make(map[string]*OccClient)
65+
var clientsForContainers sync.Map
6366

6467
const taskFinalizer string = "aliecs.alice.cern/finalizer"
6568

@@ -136,7 +139,8 @@ func (r *TaskReconciler) Reconcile(ctx context.Context, req ctrl.Request) (ctrl.
136139
return ctrl.Result{}, nil
137140
}
138141

139-
if _, exists := clientsForContainers[t.Name]; !exists {
142+
// if _, exists := clientsForContainers[t.Name]; !exists {
143+
if _, exists := clientsForContainers.Load(t.Name); !exists {
140144
if existingPod.Status.PodIP == "" {
141145
log.Info("pod doesn't have IP yet, we wait for different event")
142146
return ctrl.Result{}, nil
@@ -159,12 +163,13 @@ func (r *TaskReconciler) Reconcile(ctx context.Context, req ctrl.Request) (ctrl.
159163
// on them being implemented
160164
if t.Status.State == "" {
161165
log.V(1).Info("Status.State is empty, querying container")
162-
client, exists := clientsForContainers[t.Name]
166+
// client, exists := clientsForContainers[t.Name]
167+
client, exists := clientsForContainers.Load(t.Name)
163168
if !exists {
164169
return ctrl.Result{Requeue: true}, nil
165170
}
166171

167-
stateReply, err := client.GetState(ctx)
172+
stateReply, err := client.(*OccClient).GetState(ctx)
168173
if err != nil {
169174
log.Error(err, "Failed to GetState")
170175
return ctrl.Result{RequeueAfter: 5 * time.Second}, nil
@@ -188,16 +193,18 @@ func (r *TaskReconciler) Reconcile(ctx context.Context, req ctrl.Request) (ctrl.
188193

189194
// Handle Spec -> gRPC State Sync
190195
if t.Status.State != t.Spec.State {
191-
client, exists := clientsForContainers[t.Name]
196+
// client, exists := clientsForContainers[t.Name]
197+
client, exists := clientsForContainers.Load(t.Name)
192198
if !exists {
193199
return ctrl.Result{Requeue: true}, nil
194200
}
195201

196-
stateReply, err := client.GetState(ctx)
202+
stateReply, err := client.(*OccClient).GetState(ctx)
197203
if err != nil {
198204
log.Info("Failed to get state for sync, retrying in 5s", "error", err.Error())
199-
client.Close()
200-
delete(clientsForContainers, t.Name)
205+
client.(*OccClient).Close()
206+
// delete(clientsForContainers, t.Name)
207+
clientsForContainers.Delete(t.Name)
201208
return ctrl.Result{RequeueAfter: 5 * time.Second}, nil
202209
}
203210

@@ -208,9 +215,9 @@ func (r *TaskReconciler) Reconcile(ctx context.Context, req ctrl.Request) (ctrl.
208215
)
209216

210217
if t.Spec.Control.Mode == "fairmq" {
211-
newState, transErr = client.FairMQTransitionRequest(ctx, stateReply.GetState(), t.Spec.State, t.Spec.Arguments)
218+
newState, transErr = client.(*OccClient).FairMQTransitionRequest(ctx, stateReply.GetState(), t.Spec.State, t.Spec.Arguments)
212219
} else {
213-
reply, err := client.TransitionRequest(ctx, stateReply.GetState(), t.Spec.State, t.Spec.Arguments)
220+
reply, err := client.(*OccClient).TransitionRequest(ctx, stateReply.GetState(), t.Spec.State, t.Spec.Arguments)
214221
transErr = err
215222
if err == nil && reply.GetOk() {
216223
newState = strings.ToLower(reply.GetState())
@@ -268,7 +275,8 @@ func (r *TaskReconciler) createGRPCConsumer(ctx context.Context, t *aliecsv1alph
268275
return ctrl.Result{RequeueAfter: 5 * time.Second}, nil
269276
}
270277

271-
clientsForContainers[t.Name] = client
278+
// clientsForContainers[t.Name] = client
279+
clientsForContainers.Store(t.Name, client)
272280

273281
if err := r.recordCondition(ctx, t, aliecsv1alpha1.ConditionGRPCConnected, metav1.ConditionTrue, "Connected", fmt.Sprintf("gRPC connection established to %s", addr)); err != nil {
274282
return ctrl.Result{}, err
@@ -277,14 +285,15 @@ func (r *TaskReconciler) createGRPCConsumer(ctx context.Context, t *aliecsv1alph
277285
}
278286

279287
func (r *TaskReconciler) consumeGRPCConsumerIfReady(ctx context.Context, t *aliecsv1alpha1.Task, log logr.Logger) ctrl.Result {
280-
client, exists := clientsForContainers[t.Name]
288+
// client, exists := clientsForContainers[t.Name]
289+
client, exists := clientsForContainers.Load(t.Name)
281290

282291
if !exists {
283292
log.Info("didn't found existing client, retrying ", "task", t.Name)
284293
return ctrl.Result{RequeueAfter: time.Second}
285294
}
286295

287-
if !client.ConsumeIfReady(ctx) {
296+
if !client.(*OccClient).ConsumeIfReady(ctx) {
288297
log.Info("gRPC client is not ready, retrying in 5 seconds", "name", t.Name)
289298
return ctrl.Result{RequeueAfter: 5 * time.Second}
290299
}
@@ -342,12 +351,14 @@ func (r *TaskReconciler) deletePod(ctx context.Context, t *aliecsv1alpha1.Task,
342351
}
343352

344353
func (*TaskReconciler) cleargRPC(t *aliecsv1alpha1.Task, log logr.Logger) {
345-
if client, exists := clientsForContainers[t.Name]; exists {
354+
// if client, exists := clientsForContainers[t.Name]; exists {
355+
if client, exists := clientsForContainers.Load(t.Name); exists {
346356
log.Info("Cleaning up gRPC connection")
347-
if err := client.Close(); err != nil {
357+
if err := client.(*OccClient).Close(); err != nil {
348358
log.Error(err, "Failed to close gRPC client during deletion")
349359
}
350-
delete(clientsForContainers, t.Name)
360+
// delete(clientsForContainers, t.Name)
361+
clientsForContainers.Delete(t.Name)
351362
log.Info("gRPC cleaned")
352363
}
353364
}
@@ -432,7 +443,7 @@ func (r *TaskReconciler) SetupWithManager(mgr ctrl.Manager) error {
432443
}),
433444
)).
434445
Owns(&v1.Pod{}).
435-
WithOptions(controller.Options{MaxConcurrentReconciles: 1}).
446+
WithOptions(controller.Options{MaxConcurrentReconciles: r.MaxConcurrentReconciles}).
436447
Complete(r)
437448
}
438449

0 commit comments

Comments
 (0)