Skip to content

Commit c856c5e

Browse files
committed
[task-controller] reconcile happy path tests from Task creation until working transition
1 parent 746234a commit c856c5e

2 files changed

Lines changed: 326 additions & 41 deletions

File tree

Lines changed: 137 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,137 @@
1+
/*
2+
* === This file is part of ALICE O² ===
3+
*
4+
* Copyright 2026 CERN and copyright holders of ALICE O².
5+
* Author: Michal Tichak <michal.tichak@cern.ch>
6+
*
7+
* This program is free software: you can redistribute it and/or modify
8+
* it under the terms of the GNU General Public License as published by
9+
* the Free Software Foundation, either version 3 of the License, or
10+
* (at your option) any later version.
11+
*
12+
* This program is distributed in the hope that it will be useful,
13+
* but WITHOUT ANY WARRANTY; without even the implied warranty of
14+
* MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
15+
* GNU General Public License for more details.
16+
*
17+
* You should have received a copy of the GNU General Public License
18+
* along with this program. If not, see <http://www.gnu.org/licenses/>.
19+
*
20+
* In applying this license CERN does not waive the privileges and
21+
* immunities granted to it by virtue of its status as an
22+
* Intergovernmental Organization or submit itself to any jurisdiction.
23+
*/
24+
25+
package controller
26+
27+
import (
28+
"context"
29+
"errors"
30+
"net"
31+
"strings"
32+
"sync"
33+
34+
"google.golang.org/grpc"
35+
36+
pb "github.com/AliceO2Group/Control/control-operator/internal/controller/protos/generated"
37+
)
38+
39+
// occServerStub is an in-process implementation of the OCC gRPC service defined
40+
// in occ/protos/occ.proto. Embedding pb.UnimplementedOccServer satisfies the
41+
// generated interface and keeps this type valid if the service gains methods.
42+
// The state is guarded because the gRPC server answers on its own goroutines while
43+
// the test reads and writes it.
44+
type occServerStub struct {
45+
pb.UnimplementedOccServer
46+
47+
returnError bool
48+
mu sync.RWMutex
49+
state string
50+
}
51+
52+
func (s *occServerStub) setState(state string) {
53+
s.mu.Lock()
54+
defer s.mu.Unlock()
55+
s.state = state
56+
}
57+
58+
func (s *occServerStub) getState() string {
59+
s.mu.RLock()
60+
defer s.mu.RUnlock()
61+
return s.state
62+
}
63+
64+
func (s *occServerStub) GetState(ctx context.Context, req *pb.GetStateRequest) (*pb.GetStateReply, error) {
65+
if s.returnError == true {
66+
return nil, errors.New("returning error as requested by user")
67+
}
68+
return &pb.GetStateReply{State: s.getState()}, nil
69+
}
70+
71+
// Transition answers with the state the requested event leads to, looked up in the
72+
// same FSM table the controller used to pick the event, and moves the stub to that
73+
// state so a following GetState reports it.
74+
func (s *occServerStub) Transition(ctx context.Context, req *pb.TransitionRequest) (*pb.TransitionReply, error) {
75+
if s.returnError == true {
76+
return nil, errors.New("returning error as requested by user")
77+
}
78+
79+
src, err := StateFromString(strings.ToLower(req.GetSrcState()))
80+
if err != nil {
81+
return nil, err
82+
}
83+
event, err := TransitionFromString(strings.ToLower(req.GetTransitionEvent()))
84+
if err != nil {
85+
return nil, err
86+
}
87+
88+
for fromTo, transition := range fromStatesToTransition {
89+
if fromTo.from == src && transition == event {
90+
newState := fromTo.to.String()
91+
s.setState(newState)
92+
return &pb.TransitionReply{
93+
Trigger: pb.StateChangeTrigger_EXECUTOR,
94+
State: newState,
95+
TransitionEvent: req.GetTransitionEvent(),
96+
Ok: true,
97+
}, nil
98+
}
99+
}
100+
101+
// No rule for this event in this state, so the device stays where it is.
102+
return &pb.TransitionReply{
103+
Trigger: pb.StateChangeTrigger_EXECUTOR,
104+
State: src.String(),
105+
TransitionEvent: req.GetTransitionEvent(),
106+
Ok: false,
107+
}, nil
108+
}
109+
110+
func (s *occServerStub) StateStream(req *pb.StateStreamRequest, srv grpc.ServerStreamingServer[pb.StateStreamReply]) error {
111+
<-srv.Context().Done()
112+
return nil
113+
}
114+
115+
func (s *occServerStub) EventStream(req *pb.EventStreamRequest, srv grpc.ServerStreamingServer[pb.EventStreamReply]) error {
116+
<-srv.Context().Done()
117+
return nil
118+
}
119+
120+
// startOccServerStub serves a stub on a kernel-assigned loopback port and
121+
// returns it together with that port and a stop function.
122+
func startOccServerStub() (*occServerStub, int, func(), error) {
123+
lis, err := net.Listen("tcp", "127.0.0.1:0")
124+
if err != nil {
125+
return nil, 0, nil, err
126+
}
127+
128+
stub := &occServerStub{}
129+
srv := grpc.NewServer()
130+
pb.RegisterOccServer(srv, stub)
131+
132+
go func() {
133+
_ = srv.Serve(lis)
134+
}()
135+
136+
return stub, lis.Addr().(*net.TCPAddr).Port, srv.Stop, nil
137+
}

‎control-operator/internal/controller/task_controller_test.go‎

Lines changed: 189 additions & 41 deletions
Original file line numberDiff line numberDiff line change
@@ -26,71 +26,219 @@ package controller
2626

2727
import (
2828
"context"
29+
"time"
2930

3031
. "github.com/onsi/ginkgo/v2"
3132
. "github.com/onsi/gomega"
3233
v1 "k8s.io/api/core/v1"
3334
"k8s.io/apimachinery/pkg/api/errors"
35+
"k8s.io/apimachinery/pkg/api/meta"
36+
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
3437
"k8s.io/apimachinery/pkg/types"
38+
"k8s.io/client-go/tools/record"
39+
"sigs.k8s.io/controller-runtime/pkg/client"
40+
"sigs.k8s.io/controller-runtime/pkg/controller/controllerutil"
3541
"sigs.k8s.io/controller-runtime/pkg/reconcile"
3642

37-
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
38-
3943
aliecsv1alpha1 "github.com/AliceO2Group/Control/control-operator/api/v1alpha1"
4044
)
4145

4246
var _ = Describe("Task Controller", func() {
43-
Context("When reconciling a resource", func() {
44-
const resourceName = "test-resource"
47+
// Reconcile takes the Task one step further on every pass and then returns to wait
48+
// for the next event: first the finalizer, then the Pod, then the gRPC connection,
49+
// then the state of the container, and finally the transition. These specs run in
50+
// order and each one drives the pass responsible for the next step.
51+
Describe("When reconciling a Task step by step", Ordered, func() {
52+
const (
53+
resourceName = "test-resource"
54+
namespace = "default"
55+
containerName = "placeholder-name"
56+
containerImg = "placeholder-image:latest"
57+
// The OCC stub listens on loopback, so the Pod IP the controller dials has
58+
// to point back at the test process.
59+
podIP = "127.0.0.1"
60+
)
4561

46-
ctx := context.Background()
62+
var (
63+
ctx context.Context
64+
taskName types.NamespacedName
65+
podName types.NamespacedName
66+
occServer *occServerStub
67+
reconciler *TaskReconciler
68+
)
4769

48-
typeNamespacedName := types.NamespacedName{
49-
Name: resourceName,
50-
Namespace: "default", // TODO(user):Modify as needed
70+
reconcileOnce := func() reconcile.Result {
71+
GinkgoHelper()
72+
res, err := reconciler.Reconcile(ctx, reconcile.Request{NamespacedName: taskName})
73+
Expect(err).NotTo(HaveOccurred())
74+
return res
5175
}
52-
task := &aliecsv1alpha1.Task{}
5376

54-
BeforeEach(func() {
55-
By("creating the custom resource for the Kind Task")
56-
err := k8sClient.Get(ctx, typeNamespacedName, task)
57-
if err != nil && errors.IsNotFound(err) {
58-
resource := &aliecsv1alpha1.Task{
59-
ObjectMeta: metav1.ObjectMeta{
60-
Name: resourceName,
61-
Namespace: "default",
62-
},
63-
Spec: aliecsv1alpha1.TaskSpec{
64-
State: "standby",
65-
Pod: v1.PodSpec{Containers: []v1.Container{}},
66-
},
67-
}
68-
Expect(k8sClient.Create(ctx, resource)).To(Succeed())
69-
}
70-
})
71-
72-
AfterEach(func() {
73-
// TODO(user): Cleanup logic after each test, like removing the task instance.
77+
getTask := func() *aliecsv1alpha1.Task {
78+
GinkgoHelper()
7479
task := &aliecsv1alpha1.Task{}
75-
err := k8sClient.Get(ctx, typeNamespacedName, task)
80+
Expect(k8sClient.Get(ctx, taskName, task)).To(Succeed())
81+
return task
82+
}
83+
84+
BeforeAll(func() {
85+
ctx = context.Background()
86+
taskName = types.NamespacedName{Name: resourceName, Namespace: namespace}
87+
podName = types.NamespacedName{Name: podNameFromTask(resourceName), Namespace: namespace}
88+
89+
By("starting the OCC server stub that stands in for the task container")
90+
stub, port, stopOccServer, err := startOccServerStub()
7691
Expect(err).NotTo(HaveOccurred())
92+
DeferCleanup(stopOccServer)
7793

78-
By("Cleanup the specific resource instance Task")
79-
Expect(k8sClient.Delete(ctx, task)).To(Succeed())
80-
})
81-
It("should successfully reconcile the resource", func() {
82-
By("Reconciling the created resource")
83-
controllerReconciler := &TaskReconciler{
94+
occServer = stub
95+
occServer.setState("standby")
96+
97+
reconciler = &TaskReconciler{
8498
Client: k8sClient,
8599
Scheme: k8sClient.Scheme(),
100+
// A zero-value FakeRecorder discards events; NewFakeRecorder would block
101+
// once its buffer filled up.
102+
Recorder: &record.FakeRecorder{},
86103
}
87104

88-
_, err := controllerReconciler.Reconcile(ctx, reconcile.Request{
89-
NamespacedName: typeNamespacedName,
105+
By("creating the custom resource for the Kind Task")
106+
task := &aliecsv1alpha1.Task{
107+
ObjectMeta: metav1.ObjectMeta{
108+
Name: resourceName,
109+
Namespace: namespace,
110+
},
111+
Spec: aliecsv1alpha1.TaskSpec{
112+
Control: aliecsv1alpha1.TaskSpecControl{Mode: "direct", Port: port},
113+
State: "standby",
114+
Pod: v1.PodSpec{
115+
Containers: []v1.Container{{Name: containerName, Image: containerImg}},
116+
},
117+
},
118+
}
119+
Expect(k8sClient.Create(ctx, task)).To(Succeed())
120+
121+
// The last spec deletes the Task and its Pod through the controller, so this
122+
// only clears what an earlier failure left behind.
123+
DeferCleanup(func() {
124+
if c, loaded := clientsForContainers.LoadAndDelete(resourceName); loaded {
125+
Expect(c.(*OccClient).Close()).To(Succeed())
126+
}
127+
128+
pod := &v1.Pod{}
129+
if err := k8sClient.Get(ctx, podName, pod); err == nil {
130+
Expect(client.IgnoreNotFound(k8sClient.Delete(ctx, pod, client.GracePeriodSeconds(0)))).To(Succeed())
131+
}
132+
133+
leftover := &aliecsv1alpha1.Task{}
134+
if err := k8sClient.Get(ctx, taskName, leftover); err == nil {
135+
controllerutil.RemoveFinalizer(leftover, taskFinalizer)
136+
Expect(k8sClient.Update(ctx, leftover)).To(Succeed())
137+
Expect(client.IgnoreNotFound(k8sClient.Delete(ctx, leftover))).To(Succeed())
138+
}
90139
})
91-
Expect(err).NotTo(HaveOccurred())
92-
// TODO(user): Add more specific assertions depending on your controller's reconciliation logic.
93-
// Example: If you expect a certain status condition after reconciliation, verify it here.
140+
})
141+
142+
It("sets the finalizer before touching anything else", func() {
143+
reconcileOnce()
144+
145+
Expect(controllerutil.ContainsFinalizer(getTask(), taskFinalizer)).To(BeTrue())
146+
Expect(k8sClient.Get(ctx, podName, &v1.Pod{})).To(Satisfy(errors.IsNotFound),
147+
"the Pod should wait until the finalizer is persisted")
148+
})
149+
150+
It("creates the Pod for the Task", func() {
151+
reconcileOnce()
152+
153+
pod := &v1.Pod{}
154+
Expect(k8sClient.Get(ctx, podName, pod)).To(Succeed())
155+
156+
Expect(pod.Spec.Containers).To(HaveLen(1))
157+
Expect(pod.Spec.Containers[0].Image).To(Equal(containerImg))
158+
Expect(pod.Spec.RestartPolicy).To(Equal(v1.RestartPolicyNever))
159+
Expect(pod.Labels).To(HaveKeyWithValue("task_name", resourceName))
160+
161+
Expect(pod.OwnerReferences).To(HaveLen(1))
162+
Expect(pod.OwnerReferences[0].Name).To(Equal(resourceName))
163+
Expect(pod.OwnerReferences[0].Controller).To(HaveValue(BeTrue()))
164+
})
165+
166+
It("waits for the Pod to report an IP before connecting", func() {
167+
reconcileOnce()
168+
169+
_, connected := clientsForContainers.Load(resourceName)
170+
Expect(connected).To(BeFalse())
171+
Expect(meta.FindStatusCondition(getTask().Status.Conditions, aliecsv1alpha1.ConditionPodReady)).To(BeNil())
172+
})
173+
174+
It("connects to the container once the Pod is running", func() {
175+
By("reporting the Pod status that a kubelet would report")
176+
// envtest runs no kubelet, so nothing else moves the Pod to Running.
177+
pod := &v1.Pod{}
178+
Expect(k8sClient.Get(ctx, podName, pod)).To(Succeed())
179+
pod.Status.Phase = v1.PodRunning
180+
pod.Status.PodIP = podIP
181+
Expect(k8sClient.Status().Update(ctx, pod)).To(Succeed())
182+
183+
reconcileOnce()
184+
185+
occClient, connected := clientsForContainers.Load(resourceName)
186+
Expect(connected).To(BeTrue())
187+
188+
task := getTask()
189+
Expect(meta.IsStatusConditionTrue(task.Status.Conditions, aliecsv1alpha1.ConditionPodReady)).To(BeTrue())
190+
Expect(meta.IsStatusConditionTrue(task.Status.Conditions, aliecsv1alpha1.ConditionGRPCConnected)).To(BeTrue())
191+
192+
// This pass asked to be requeued until the connection is Ready, which is what
193+
// the next step needs.
194+
connectCtx, cancel := context.WithTimeout(ctx, 5*time.Second)
195+
defer cancel()
196+
Expect(occClient.(*OccClient).WaitUntilConnected(connectCtx)).To(Succeed())
197+
})
198+
199+
It("reads the state from the container", func() {
200+
reconcileOnce()
201+
202+
task := getTask()
203+
Expect(task.Status.State).To(Equal(occServer.getState()))
204+
Expect(task.Status.Error).To(BeEmpty())
205+
Expect(meta.IsStatusConditionTrue(task.Status.Conditions, aliecsv1alpha1.ConditionStateAccessible)).To(BeTrue())
206+
})
207+
208+
It("transitions the container to a newly requested state", func() {
209+
By("asking the Task for a state the container is not in")
210+
task := getTask()
211+
task.Spec.State = "configured"
212+
Expect(k8sClient.Update(ctx, task)).To(Succeed())
213+
214+
reconcileOnce()
215+
216+
task = getTask()
217+
Expect(task.Status.State).To(Equal("configured"))
218+
Expect(task.Status.Error).To(BeEmpty())
219+
Expect(meta.IsStatusConditionTrue(task.Status.Conditions, aliecsv1alpha1.ConditionStateTransitionResult)).To(BeTrue())
220+
221+
By("leaving the container in the requested state")
222+
Expect(occServer.getState()).To(Equal("configured"))
223+
})
224+
225+
It("clears the connection and the Pod when the Task is deleted", func() {
226+
Expect(k8sClient.Delete(ctx, getTask())).To(Succeed())
227+
228+
reconcileOnce()
229+
230+
_, connected := clientsForContainers.Load(resourceName)
231+
Expect(connected).To(BeFalse())
232+
233+
// Nothing schedules the Pod in envtest, so the API server drops it right away
234+
// instead of waiting for a kubelet to confirm termination.
235+
Expect(k8sClient.Get(ctx, podName, &v1.Pod{})).To(Satisfy(errors.IsNotFound))
236+
Expect(controllerutil.ContainsFinalizer(getTask(), taskFinalizer)).To(BeTrue(),
237+
"this pass only deletes the Pod, the finalizer goes on the next one")
238+
239+
reconcileOnce()
240+
241+
Expect(k8sClient.Get(ctx, taskName, &aliecsv1alpha1.Task{})).To(Satisfy(errors.IsNotFound))
94242
})
95243
})
96244
})

0 commit comments

Comments
 (0)