mirror of
https://github.com/wahyd4/kt-connect.git
synced 2026-08-09 05:16:02 +10:00
refact: make cloud watcher stopable
This commit is contained in:
@@ -4,7 +4,6 @@ import (
|
||||
"time"
|
||||
|
||||
"k8s.io/apimachinery/pkg/util/runtime"
|
||||
"k8s.io/apimachinery/pkg/util/wait"
|
||||
"k8s.io/client-go/informers"
|
||||
v1 "k8s.io/client-go/listers/core/v1"
|
||||
"k8s.io/client-go/tools/cache"
|
||||
@@ -20,10 +19,9 @@ func podDeleted(obj interface{}) {
|
||||
}
|
||||
|
||||
// Pods watch pods change
|
||||
func (w *Watcher) Pods() (lister v1.PodLister, err error) {
|
||||
func (w *Watcher) Pods(stopCh <-chan struct{}) (lister v1.PodLister, err error) {
|
||||
resyncPeriod := 30 * time.Minute
|
||||
|
||||
stopCh := wait.NeverStop
|
||||
factory := informers.NewSharedInformerFactory(w.Client, resyncPeriod)
|
||||
podInformer := factory.Core().V1().Pods()
|
||||
informer := podInformer.Informer()
|
||||
|
||||
@@ -34,7 +34,7 @@ func Construct(client kubernetes.Interface, config *rest.Config) (w Watcher, err
|
||||
}
|
||||
w.NamespaceLister = namespaceLister
|
||||
|
||||
podListener, err := w.Pods()
|
||||
podListener, err := w.Pods(wait.NeverStop)
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
@@ -55,9 +55,9 @@ func Construct(client kubernetes.Interface, config *rest.Config) (w Watcher, err
|
||||
}
|
||||
|
||||
// ServiceListener ServiceListener
|
||||
func ServiceListener(client kubernetes.Interface) (lister v1.ServiceLister, err error) {
|
||||
func ServiceListener(client kubernetes.Interface, stopCh <-chan struct{}) (lister v1.ServiceLister, err error) {
|
||||
w := Watcher{Client: client}
|
||||
lister, err = w.Services(wait.NeverStop)
|
||||
lister, err = w.Services(stopCh)
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
@@ -65,9 +65,9 @@ func ServiceListener(client kubernetes.Interface) (lister v1.ServiceLister, err
|
||||
}
|
||||
|
||||
// PodListener PodListener
|
||||
func PodListener(client kubernetes.Interface) (lister v1.PodLister, err error) {
|
||||
func PodListener(client kubernetes.Interface, stopCh <-chan struct{}) (lister v1.PodLister, err error) {
|
||||
w := Watcher{Client: client}
|
||||
lister, err = w.Pods()
|
||||
lister, err = w.Pods(stopCh)
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
|
||||
@@ -3,6 +3,7 @@ package cluster
|
||||
import (
|
||||
clusterWatcher "github.com/alibaba/kt-connect/pkg/apiserver/cluster"
|
||||
appV1 "k8s.io/api/apps/v1"
|
||||
"k8s.io/apimachinery/pkg/util/wait"
|
||||
"k8s.io/client-go/kubernetes"
|
||||
v1 "k8s.io/client-go/listers/core/v1"
|
||||
)
|
||||
@@ -13,8 +14,8 @@ func Create(kubeConfig string) (kubernetes Kubernetes, err error) {
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
serviceListener, err := clusterWatcher.ServiceListener(clientSet)
|
||||
podListener, err := clusterWatcher.PodListener(clientSet)
|
||||
serviceListener, err := clusterWatcher.ServiceListener(clientSet, wait.NeverStop)
|
||||
podListener, err := clusterWatcher.PodListener(clientSet, wait.NeverStop)
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user