-
Notifications
You must be signed in to change notification settings - Fork 1
/
depending_on_pipeline_reconciler.go
111 lines (91 loc) · 3.29 KB
/
depending_on_pipeline_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
package pipelines
import (
"context"
"github.com/sky-uk/kfp-operator/apis"
pipelinesv1 "github.com/sky-uk/kfp-operator/apis/pipelines/v1alpha5"
"k8s.io/apimachinery/pkg/api/errors"
"k8s.io/apimachinery/pkg/types"
ctrl "sigs.k8s.io/controller-runtime"
"sigs.k8s.io/controller-runtime/pkg/builder"
"sigs.k8s.io/controller-runtime/pkg/client"
"sigs.k8s.io/controller-runtime/pkg/handler"
"sigs.k8s.io/controller-runtime/pkg/log"
"sigs.k8s.io/controller-runtime/pkg/predicate"
"sigs.k8s.io/controller-runtime/pkg/reconcile"
"sigs.k8s.io/controller-runtime/pkg/source"
)
const (
pipelineRefField = ".spec.pipeline"
)
type DependingOnPipelineResource interface {
client.Object
GetPipeline() pipelinesv1.PipelineIdentifier
GetObservedPipelineVersion() string
SetObservedPipelineVersion(string)
}
type DependingOnPipelineReconciler[R DependingOnPipelineResource] struct {
EC K8sExecutionContext
}
func (dr DependingOnPipelineReconciler[R]) handleObservedPipelineVersion(ctx context.Context, pipelineIdentifier pipelinesv1.PipelineIdentifier, resource R) (bool, error) {
logger := log.FromContext(ctx)
setVersion := true
desiredVersion := pipelineIdentifier.Version
if pipelineIdentifier.Version == "" {
pipeline, err := dr.getIgnoreNotFound(ctx, types.NamespacedName{
Namespace: resource.GetNamespace(),
Name: pipelineIdentifier.Name,
})
if err != nil {
return false, err
}
desiredVersion, setVersion = dependentPipelineVersionIfSucceeded(pipeline)
}
if setVersion && resource.GetObservedPipelineVersion() != desiredVersion {
resource.SetObservedPipelineVersion(desiredVersion)
if err := dr.EC.Client.Status().Update(ctx, resource); err != nil {
logger.Error(err, "error updating resource with observed pipeline version")
return false, err
}
return true, nil
}
return false, nil
}
func (dr DependingOnPipelineReconciler[R]) getIgnoreNotFound(ctx context.Context, key client.ObjectKey) (*pipelinesv1.Pipeline, error) {
logger := log.FromContext(ctx)
pipeline := &pipelinesv1.Pipeline{}
if err := dr.EC.Client.NonCached.Get(ctx, key, pipeline); err != nil {
if errors.IsNotFound(err) {
logger.Info("object not found")
return nil, nil
}
logger.Error(err, "unable to fetch object")
return nil, err
}
return pipeline, nil
}
func (dr DependingOnPipelineReconciler[R]) setupWithManager(mgr ctrl.Manager, controllerBuilder *builder.Builder, object client.Object, reconciliationRequestsForPipeline func(client.Object) []reconcile.Request) (*builder.Builder, error) {
if err := mgr.GetFieldIndexer().IndexField(context.Background(), object, pipelineRefField, func(rawObj client.Object) []string {
referencingResource := rawObj.(R)
return []string{referencingResource.GetPipeline().Name}
}); err != nil {
return nil, err
}
return controllerBuilder.Watches(
&source.Kind{Type: &pipelinesv1.Pipeline{}},
handler.EnqueueRequestsFromMapFunc(reconciliationRequestsForPipeline),
builder.WithPredicates(predicate.ResourceVersionChangedPredicate{}),
), nil
}
func dependentPipelineVersionIfSucceeded(pipeline *pipelinesv1.Pipeline) (string, bool) {
if pipeline == nil {
return "", true
}
switch pipeline.Status.SynchronizationState {
case apis.Succeeded:
return pipeline.Status.Version, true
case apis.Deleted:
return "", true
default:
return "", false
}
}