forked from kubeflow/common
/
test_job_reconciler.go
163 lines (124 loc) · 4.33 KB
/
test_job_reconciler.go
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
package test_job
import (
"context"
commonv1 "github.com/dafu-wu/common/pkg/apis/common/v1"
common_reconciler "github.com/dafu-wu/common/pkg/reconciler.v1/common"
v1 "github.com/dafu-wu/common/test_job/apis/test_job/v1"
"github.com/dafu-wu/common/test_job/client/clientset/versioned/scheme"
corev1 "k8s.io/api/core/v1"
"k8s.io/apimachinery/pkg/runtime"
utilruntime "k8s.io/apimachinery/pkg/util/runtime"
clientgoscheme "k8s.io/client-go/kubernetes/scheme"
ctrl "sigs.k8s.io/controller-runtime"
"sigs.k8s.io/controller-runtime/pkg/client"
"sigs.k8s.io/controller-runtime/pkg/log"
)
type TestReconciler struct {
common_reconciler.ReconcilerUtil
common_reconciler.ServiceReconciler
common_reconciler.PodReconciler
common_reconciler.VolcanoReconciler
common_reconciler.JobReconciler
DC *DummyClient
Job *v1.TestJob
Pods []*corev1.Pod
Services []*corev1.Service
PodGroup client.Object
}
func NewTestReconciler() *TestReconciler {
scheme := runtime.NewScheme()
utilruntime.Must(clientgoscheme.AddToScheme(scheme))
utilruntime.Must(v1.AddToScheme(scheme))
dummy_client := &DummyClient{}
r := &TestReconciler{
DC: dummy_client,
}
// Generate Bare Components
jobR := common_reconciler.BareJobReconciler(dummy_client)
jobR.OverrideForJobInterface(r, r, r, r)
podR := common_reconciler.BarePodReconciler(dummy_client)
podR.OverrideForPodInterface(r, r, r)
svcR := common_reconciler.BareServiceReconciler(dummy_client)
svcR.OverrideForServiceInterface(r, r, r)
gangR := common_reconciler.BareVolcanoReconciler(dummy_client, nil, false)
gangR.OverrideForGangSchedulingInterface(r)
Log := log.Log
utilR := common_reconciler.BareUtilReconciler(nil, Log, scheme)
//kubeflowReconciler := common_reconciler.BareKubeflowReconciler()
r.JobReconciler = *jobR
r.PodReconciler = *podR
r.ServiceReconciler = *svcR
r.VolcanoReconciler = *gangR
r.ReconcilerUtil = *utilR
return r
}
// Reconcile is part of the main kubernetes reconciliation loop which aims to
// move the current state of the cluster closer to the desired state.
func (r *TestReconciler) Reconcile(ctx context.Context, req ctrl.Request) (ctrl.Result, error) {
_ = log.FromContext(ctx)
job, err := r.GetJob(ctx, req)
if err != nil {
return ctrl.Result{}, err
}
logger := r.GetLogger(job)
if job.GetDeletionTimestamp() != nil {
return ctrl.Result{}, nil
}
scheme.Scheme.Default(job)
// Get rid of SatisfiedExpectation
replicasSpec, err := r.ExtractReplicasSpec(job)
if err != nil {
return ctrl.Result{}, err
}
runPolicy, err := r.ExtractRunPolicy(job)
if err != nil {
return ctrl.Result{}, err
}
status, err := r.ExtractJobStatus(job)
if err != nil {
return ctrl.Result{}, err
}
err = r.ReconcileJob(ctx, job, replicasSpec, status, runPolicy)
if err != nil {
logger.Info("Reconcile Test Job error %v", err)
return ctrl.Result{}, err
}
return ctrl.Result{}, nil
}
func (r *TestReconciler) GetReconcilerName() string {
return "Test Reconciler"
}
func (r *TestReconciler) GetJob(ctx context.Context, req ctrl.Request) (client.Object, error) {
return r.Job, nil
}
func (r *TestReconciler) GetDefaultContainerName() string {
return v1.DefaultContainerName
}
func (r *TestReconciler) GetPodGroupForJob(ctx context.Context, job client.Object) (client.Object, error) {
return r.PodGroup, nil
}
func (r *TestReconciler) GetPodsForJob(ctx context.Context, job client.Object) ([]*corev1.Pod, error) {
return r.Pods, nil
}
func (r *TestReconciler) GetServicesForJob(ctx context.Context, job client.Object) ([]*corev1.Service, error) {
return r.Services, nil
}
func (r *TestReconciler) ExtractReplicasSpec(job client.Object) (map[commonv1.ReplicaType]*commonv1.ReplicaSpec, error) {
tj := job.(*v1.TestJob)
rs := map[commonv1.ReplicaType]*commonv1.ReplicaSpec{}
for k, v := range tj.Spec.TestReplicaSpecs {
rs[commonv1.ReplicaType(k)] = v
}
return rs, nil
}
func (r *TestReconciler) ExtractRunPolicy(job client.Object) (*commonv1.RunPolicy, error) {
tj := job.(*v1.TestJob)
return tj.Spec.RunPolicy, nil
}
func (r *TestReconciler) ExtractJobStatus(job client.Object) (*commonv1.JobStatus, error) {
tj := job.(*v1.TestJob)
return &tj.Status, nil
}
func (r *TestReconciler) IsMasterRole(replicas map[commonv1.ReplicaType]*commonv1.ReplicaSpec, rtype commonv1.ReplicaType, index int) bool {
return string(rtype) == string(v1.TestReplicaTypeMaster)
}