diff --git a/pkg/fission-cli/cmd/kubewatch/command.go b/pkg/fission-cli/cmd/kubewatch/command.go index 7c3eef58..23e38ea0 100644 --- a/pkg/fission-cli/cmd/kubewatch/command.go +++ b/pkg/fission-cli/cmd/kubewatch/command.go @@ -43,8 +43,8 @@ func Commands() *cobra.Command { RunE: wrapper.Wrapper(Delete), } wrapper.SetFlags(deleteCmd, flag.FlagSet{ - Required: []flag.Flag{flag.KwFnName}, - Optional: []flag.Flag{flag.NamespaceTrigger, flag.IgnoreNotFound}, + Required: []flag.Flag{flag.KwName}, + Optional: []flag.Flag{flag.NamespaceTrigger, flag.IgnoreNotFound, flag.KwFnName}, }) listCmd := &cobra.Command{ diff --git a/pkg/kubewatcher/kubewatcher.go b/pkg/kubewatcher/kubewatcher.go index eb13acf2..80af2d0e 100644 --- a/pkg/kubewatcher/kubewatcher.go +++ b/pkg/kubewatcher/kubewatcher.go @@ -42,18 +42,11 @@ import ( "github.com/fission/fission/pkg/utils" ) -type requestType int - -const ( - SYNC requestType = iota -) - type ( KubeWatcher struct { logger *zap.Logger watches map[types.UID]watchSubscription kubernetesClient kubernetes.Interface - requestChannel chan *kubeWatcherRequest publisher publisher.Publisher } @@ -66,15 +59,6 @@ type ( kubernetesClient kubernetes.Interface publisher publisher.Publisher } - - kubeWatcherRequest struct { - requestType - watches []fv1.KubernetesWatchTrigger - responseChannel chan *kubeWatcherResponse - } - kubeWatcherResponse struct { - error - } ) func MakeKubeWatcher(ctx context.Context, logger *zap.Logger, kubernetesClient kubernetes.Interface, publisher publisher.Publisher) *KubeWatcher { @@ -83,49 +67,10 @@ func MakeKubeWatcher(ctx context.Context, logger *zap.Logger, kubernetesClient k watches: make(map[types.UID]watchSubscription), kubernetesClient: kubernetesClient, publisher: publisher, - requestChannel: make(chan *kubeWatcherRequest), } - go kw.svc(ctx) return kw } -func (kw *KubeWatcher) Sync(watches []fv1.KubernetesWatchTrigger) error { - req := &kubeWatcherRequest{ - requestType: SYNC, - watches: watches, - responseChannel: make(chan *kubeWatcherResponse), - } - kw.requestChannel <- req - resp := <-req.responseChannel - return resp.error -} - -func (kw *KubeWatcher) svc(ctx context.Context) { - for { - req := <-kw.requestChannel - switch req.requestType { - case SYNC: - newWatchUids := make(map[types.UID]bool) - for _, w := range req.watches { - newWatchUids[w.ObjectMeta.UID] = true - } - // Remove old watches - for uid, ws := range kw.watches { - if _, ok := newWatchUids[uid]; !ok { - kw.removeWatch(&ws.watch) //nolint: errCheck - } - } - // Add new watches - for _, w := range req.watches { - if _, ok := kw.watches[w.ObjectMeta.UID]; !ok { - kw.addWatch(ctx, &w) //nolint: errCheck - } - } - req.responseChannel <- &kubeWatcherResponse{error: nil} - } - } -} - // TODO lifted from kubernetes/pkg/kubectl/resource_printer.go. func printKubernetesObject(obj runtime.Object, w io.Writer) error { switch obj := obj.(type) { diff --git a/pkg/kubewatcher/main.go b/pkg/kubewatcher/main.go index 6dc3afa4..4d876dc6 100644 --- a/pkg/kubewatcher/main.go +++ b/pkg/kubewatcher/main.go @@ -39,7 +39,8 @@ func Start(ctx context.Context, logger *zap.Logger, routerUrl string) error { poster := publisher.MakeWebhookPublisher(logger, routerUrl) kubeWatch := MakeKubeWatcher(ctx, logger, kubeClient, poster) - MakeWatchSync(ctx, logger, fissionClient, kubeWatch) + ws := MakeWatchSync(ctx, logger, fissionClient, kubeWatch) + ws.Run(ctx) return nil } diff --git a/pkg/kubewatcher/watchSync.go b/pkg/kubewatcher/watchSync.go index 35dfd74a..3177ade9 100644 --- a/pkg/kubewatcher/watchSync.go +++ b/pkg/kubewatcher/watchSync.go @@ -21,16 +21,19 @@ import ( "time" "go.uber.org/zap" - metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + k8sCache "k8s.io/client-go/tools/cache" + fv1 "github.com/fission/fission/pkg/apis/core/v1" "github.com/fission/fission/pkg/generated/clientset/versioned" + "github.com/fission/fission/pkg/utils" ) type ( WatchSync struct { - logger *zap.Logger - client versioned.Interface - kubeWatcher *KubeWatcher + logger *zap.Logger + client versioned.Interface + kubeWatcher *KubeWatcher + kubeWatcherInformer map[string]k8sCache.SharedIndexInformer } ) @@ -40,22 +43,28 @@ func MakeWatchSync(ctx context.Context, logger *zap.Logger, client versioned.Int client: client, kubeWatcher: kubeWatcher, } - go ws.syncSvc(ctx) + ws.kubeWatcherInformer = utils.GetInformersForNamespaces(client, time.Minute*30, fv1.KubernetesWatchResource) + ws.KubeWatcherEventHandlers(ctx) return ws } -func (ws *WatchSync) syncSvc(ctx context.Context) { - // TODO watch instead of polling - for { - watches, err := ws.client.CoreV1().KubernetesWatchTriggers(metav1.NamespaceAll).List(ctx, metav1.ListOptions{}) - if err != nil { - ws.logger.Fatal("failed to get Kubernetes watch trigger list", zap.Error(err)) - } - - err = ws.kubeWatcher.Sync(watches.Items) - if err != nil { - ws.logger.Fatal("failed to sync watches", zap.Error(err)) - } - time.Sleep(3 * time.Second) +func (ws *WatchSync) Run(ctx context.Context) { + for _, informer := range ws.kubeWatcherInformer { + go informer.Run(ctx.Done()) + } +} + +func (ws *WatchSync) KubeWatcherEventHandlers(ctx context.Context) { + for _, informer := range ws.kubeWatcherInformer { + informer.AddEventHandler(k8sCache.ResourceEventHandlerFuncs{ + AddFunc: func(obj interface{}) { + objKubeWatcher := obj.(*fv1.KubernetesWatchTrigger) + ws.kubeWatcher.addWatch(ctx, objKubeWatcher) //nolint: errCheck + }, + DeleteFunc: func(obj interface{}) { + objKubeWatcher := obj.(*fv1.KubernetesWatchTrigger) + ws.kubeWatcher.removeWatch(objKubeWatcher) //nolint: errCheck + }, + }) } }