Skip to content

Commit 4ada65d

Browse files
authored
[OCTRL-1091] Hostname resolving in environment controller (#839)
[control-operator] hostname resolving in EnvironmentReconciler Environment CRD controller needs a way how to resolve hostname passed from ECS to the one existing inside k8s cluster. ECS passes non fully qualified hostname, while kubernetes needs fully qualified hostname
1 parent 01cca13 commit 4ada65d

2 files changed

Lines changed: 57 additions & 18 deletions

File tree

control-operator/ecs-manifests/environment/test_env.yaml

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -7,6 +7,7 @@ metadata:
77
taskTemplates:
88
tasks:
99
mtichak-ost.cern.ch:
10+
# mtichak-ost:
1011
# mtichak-xps-cern:
1112
# - name: readout
1213
# argsCLI:
@@ -188,7 +189,8 @@ taskTemplates:
188189
spec:
189190
state: standby
190191
tasks:
191-
mtichak-ost.cern.ch:
192+
# mtichak-ost.cern.ch:
193+
mtichak-ost:
192194
- name: readout
193195
spec:
194196
pod:

control-operator/internal/controller/environment_controller.go

Lines changed: 54 additions & 17 deletions
Original file line numberDiff line numberDiff line change
@@ -53,7 +53,7 @@ type EnvironmentReconciler struct {
5353
}
5454

5555
func (r *EnvironmentReconciler) runTasksFromReferenceOnNode(ctx context.Context, taskReferences []aliecsv1alpha1.TaskReference,
56-
nodename string, req ctrl.Request, environment *aliecsv1alpha1.Environment, log logr.Logger,
56+
nodename string, resolvedNodename string, req ctrl.Request, environment *aliecsv1alpha1.Environment, log logr.Logger,
5757
) (*ctrl.Result, error) {
5858
for _, taskReference := range taskReferences {
5959
log.Info("geting stored template for task", "task", taskReference.Name)
@@ -106,7 +106,8 @@ func (r *EnvironmentReconciler) runTasksFromReferenceOnNode(ctx context.Context,
106106
task.Spec.Pod.Containers[0].Args = append(task.Spec.Pod.Containers[0].Args, taskReference.ArgsCLI...)
107107
maps.Copy(task.Spec.Arguments, taskReference.ArgsTransition)
108108

109-
task.Spec.Pod.NodeName = nodename
109+
task.Spec.Pod.NodeName = resolvedNodename
110+
task.Spec.NodeName = resolvedNodename
110111
task.Spec.State = environment.Spec.State
111112

112113
task.Labels = labels(environment, nodename)
@@ -131,7 +132,7 @@ func (r *EnvironmentReconciler) runTasksFromReferenceOnNode(ctx context.Context,
131132
}
132133

133134
func (r *EnvironmentReconciler) runTaskFromDefinitionOnNode(ctx context.Context, taskDefs []aliecsv1alpha1.TaskDefinition,
134-
nodename string, req ctrl.Request, environment *aliecsv1alpha1.Environment, log logr.Logger,
135+
nodename string, resolvedNodename string, req ctrl.Request, environment *aliecsv1alpha1.Environment, log logr.Logger,
135136
) (*ctrl.Result, error) {
136137
for _, taskDef := range taskDefs {
137138
task := &aliecsv1alpha1.Task{}
@@ -151,7 +152,8 @@ func (r *EnvironmentReconciler) runTaskFromDefinitionOnNode(ctx context.Context,
151152
maps.Copy(task.Spec.Arguments, taskDef.Spec.Arguments)
152153
}
153154

154-
task.Spec.Pod.NodeName = nodename
155+
task.Spec.Pod.NodeName = resolvedNodename
156+
task.Spec.NodeName = resolvedNodename
155157
task.Spec.State = environment.Spec.State
156158

157159
task.Labels = labels(environment, nodename)
@@ -249,24 +251,31 @@ func (r *EnvironmentReconciler) Reconcile(ctx context.Context, req ctrl.Request)
249251
}
250252

251253
if environment.Status.State == "" {
254+
nodeNames, err := r.listAndParseNodes(ctx)
255+
if err != nil {
256+
return ctrl.Result{}, err
257+
}
258+
252259
for nodename, tasksReferences := range environment.TaskTemplates.Tasks {
253260
log.Info("creating tasks for hostname from references", "hostname", nodename, "number of tasks", len(tasksReferences))
254-
if res, err := r.checkNodeExistence(ctx, nodename); err != nil {
255-
return res, err
261+
resolvedName, err := resolveNodeName(nodeNames, nodename)
262+
if err != nil {
263+
return ctrl.Result{}, err
256264
}
257265

258-
if res, err := r.runTasksFromReferenceOnNode(ctx, tasksReferences, nodename, req, environment, log); res != nil {
266+
if res, err := r.runTasksFromReferenceOnNode(ctx, tasksReferences, nodename, resolvedName, req, environment, log); res != nil {
259267
return *res, err
260268
}
261269
}
262270

263271
for nodename, taskTemplates := range environment.Spec.Tasks {
264272
log.Info("creating tasks for hostname from definitions", "hostname", nodename, "number of tasks", len(taskTemplates))
265-
if res, err := r.checkNodeExistence(ctx, nodename); err != nil {
266-
return res, err
273+
resolvedName, err := resolveNodeName(nodeNames, nodename)
274+
if err != nil {
275+
return ctrl.Result{}, err
267276
}
268277

269-
if res, err := r.runTaskFromDefinitionOnNode(ctx, taskTemplates, nodename, req, environment, log); res != nil {
278+
if res, err := r.runTaskFromDefinitionOnNode(ctx, taskTemplates, nodename, resolvedName, req, environment, log); res != nil {
270279
return *res, err
271280
}
272281
}
@@ -314,15 +323,43 @@ func (r *EnvironmentReconciler) Reconcile(ctx context.Context, req ctrl.Request)
314323
return ctrl.Result{}, nil
315324
}
316325

317-
func (r *EnvironmentReconciler) checkNodeExistence(ctx context.Context, nodename string) (ctrl.Result, error) {
318-
node := &v1.Node{}
319-
if err := r.Get(ctx, types.NamespacedName{Name: nodename}, node); err != nil {
320-
if k8serrors.IsNotFound(err) {
321-
return ctrl.Result{}, fmt.Errorf("node %s not found in cluster", nodename)
326+
var filterList = []string{"cern", "ch"}
327+
328+
// listAndParseNodes asks k8s cluster for all nodes and it creates map out of them containing the node names
329+
// and also partial node names while filtering those from `filterList`, so map will contain following for example:
330+
// { "flp001.cern.ch": "flp001.cern.ch", "flp001": "flp001.cern.ch" }
331+
func (r *EnvironmentReconciler) listAndParseNodes(ctx context.Context) (map[string]string, error) {
332+
nodeList := &v1.NodeList{}
333+
nodeNames := make(map[string]string)
334+
if err := r.List(ctx, nodeList); err != nil {
335+
return nodeNames, fmt.Errorf("listing nodes: %w", err)
336+
}
337+
338+
for _, node := range nodeList.Items {
339+
if val, ok := nodeNames[node.Name]; ok {
340+
return nil, fmt.Errorf("duplicate node names %s and %s ", node.Name, val)
341+
}
342+
nodeNames[node.Name] = node.Name
343+
for _, nodePart := range strings.Split(node.Name, ".") {
344+
if !slices.Contains(filterList, nodePart) {
345+
if val, ok := nodeNames[nodePart]; ok {
346+
return nil, fmt.Errorf("duplicate partial node names %s and %s ", nodePart, val)
347+
}
348+
nodeNames[nodePart] = node.Name
349+
}
322350
}
323-
return ctrl.Result{}, err
324351
}
325-
return ctrl.Result{}, nil
352+
353+
return nodeNames, nil
354+
}
355+
356+
// resolveNodeNames checks created map in listAndParseNodes for any match
357+
func resolveNodeName(nodes map[string]string, nodename string) (string, error) {
358+
if res, ok := nodes[nodename]; !ok {
359+
return "", fmt.Errorf("node not in cluster: %s", nodename)
360+
} else {
361+
return res, nil
362+
}
326363
}
327364

328365
func aggregateState(tasks []aliecsv1alpha1.Task, previousState string) string {

0 commit comments

Comments
 (0)