This repository has been archived by the owner on Jul 11, 2023. It is now read-only.
-
Notifications
You must be signed in to change notification settings - Fork 279
/
client.go
106 lines (88 loc) · 3.47 KB
/
client.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
package ingress
import (
"reflect"
extensionsV1beta "k8s.io/api/extensions/v1beta1"
"k8s.io/client-go/informers"
"k8s.io/client-go/kubernetes"
"k8s.io/client-go/tools/cache"
"github.com/openservicemesh/osm/pkg/configurator"
k8s "github.com/openservicemesh/osm/pkg/kubernetes"
"github.com/openservicemesh/osm/pkg/service"
)
// NewIngressClient implements ingress.Monitor and creates the Kubernetes client to monitor Ingress resources.
func NewIngressClient(kubeClient kubernetes.Interface, kubeController k8s.Controller, stop chan struct{}, cfg configurator.Configurator) (Monitor, error) {
informerFactory := informers.NewSharedInformerFactory(kubeClient, k8s.DefaultKubeEventResyncInterval)
informer := informerFactory.Extensions().V1beta1().Ingresses().Informer()
client := Client{
informer: informer,
cache: informer.GetStore(),
cacheSynced: make(chan interface{}),
announcements: make(chan interface{}),
kubeController: kubeController,
}
shouldObserve := func(obj interface{}) bool {
ns := reflect.ValueOf(obj).Elem().FieldByName("ObjectMeta").FieldByName("Namespace").String()
return kubeController.IsMonitoredNamespace(ns)
}
informer.AddEventHandler(k8s.GetKubernetesEventHandlers("Ingress", "Kubernetes", client.announcements, shouldObserve))
if err := client.run(stop); err != nil {
log.Error().Err(err).Msg("Could not start Kubernetes Ingress client")
return nil, err
}
return client, nil
}
// run executes informer collection.
func (c *Client) run(stop <-chan struct{}) error {
log.Info().Msg("Ingress client started")
if c.informer == nil {
return errInitInformers
}
go c.informer.Run(stop)
log.Info().Msgf("Waiting for Ingress informer cache sync")
if !cache.WaitForCacheSync(stop, c.informer.HasSynced) {
return errSyncingCaches
}
// Closing the cacheSynced channel signals to the rest of the system that... caches have been synced.
close(c.cacheSynced)
log.Info().Msgf("Cache sync finished for Ingress informer")
return nil
}
// GetAnnouncementsChannel returns the announcement channel for the Ingress client
func (c Client) GetAnnouncementsChannel() <-chan interface{} {
return c.announcements
}
// GetIngressResources returns the ingress resources whose backends correspond to the service
func (c Client) GetIngressResources(meshService service.MeshService) ([]*extensionsV1beta.Ingress, error) {
var ingressResources []*extensionsV1beta.Ingress
for _, ingressInterface := range c.cache.List() {
ingress, ok := ingressInterface.(*extensionsV1beta.Ingress)
if !ok {
log.Error().Msg("Failed type assertion for Ingress in ingress cache")
continue
}
// Extra safety - make sure we do not pay attention to Ingresses outside of observed namespaces
if !c.kubeController.IsMonitoredNamespace(ingress.Namespace) {
continue
}
// Check if the ingress resource belongs to the same namespace as the service
if ingress.Namespace != meshService.Namespace {
// The ingress resource does not belong to the namespace of the service
continue
}
if backend := ingress.Spec.Backend; backend != nil && backend.ServiceName == meshService.Name {
// Default backend service
ingressResources = append(ingressResources, ingress)
continue
}
ingressRule:
for _, rule := range ingress.Spec.Rules {
for _, path := range rule.HTTP.Paths {
if path.Backend.ServiceName == meshService.Name {
ingressResources = append(ingressResources, ingress)
break ingressRule
}
}
}
}
return ingressResources, nil
}