From 8214394af55da89c0147cd5e9ba0e42fd15b2f52 Mon Sep 17 00:00:00 2001 From: Rami Berman Date: Sun, 31 Oct 2021 14:57:49 +0200 Subject: [PATCH] WIP --- cli/cmd/tapRunner.go | 286 ++++++------------- cli/mizu.go | 2 +- shared/go.mod | 3 + shared/go.sum | 1 + {cli/mizu => shared}/goUtils/funcWrappers.go | 0 shared/kubernetes/k8sTapManager.go | 199 ++++++++++--- shared/kubernetes/k8sTapManagerErrors.go | 19 ++ shared/kubernetes/provider.go | 3 +- shared/kubernetes/utils.go | 44 ++- 9 files changed, 305 insertions(+), 252 deletions(-) rename {cli/mizu => shared}/goUtils/funcWrappers.go (100%) create mode 100644 shared/kubernetes/k8sTapManagerErrors.go diff --git a/cli/cmd/tapRunner.go b/cli/cmd/tapRunner.go index bed740029..2fb8c0934 100644 --- a/cli/cmd/tapRunner.go +++ b/cli/cmd/tapRunner.go @@ -2,9 +2,9 @@ package cmd import ( "context" - "encoding/json" "errors" "fmt" + "github.com/up9inc/mizu/shared/goUtils" "io/ioutil" k8serrors "k8s.io/apimachinery/pkg/api/errors" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" @@ -25,11 +25,9 @@ import ( "github.com/up9inc/mizu/cli/mizu" "github.com/up9inc/mizu/cli/mizu/fsUtils" - "github.com/up9inc/mizu/cli/mizu/goUtils" "github.com/up9inc/mizu/cli/telemetry" "github.com/up9inc/mizu/cli/uiUtils" "github.com/up9inc/mizu/shared" - "github.com/up9inc/mizu/shared/debounce" "github.com/up9inc/mizu/shared/kubernetes" "github.com/up9inc/mizu/shared/logger" "github.com/up9inc/mizu/tap/api" @@ -37,12 +35,11 @@ import ( const ( cleanupTimeout = time.Minute - updateTappersDelay = 5 * time.Second ) type tapState struct { apiServerService *core.Service - currentlyTappedPods []core.Pod + tapManager *kubernetes.K8sTapManager mizuServiceAccountExists bool } @@ -116,19 +113,6 @@ func RunMizuTap() { logger.Log.Infof("Tapping pods in %s", namespacesStr) - if err, _ := updateCurrentlyTappedPods(kubernetesProvider, ctx, targetNamespaces); err != nil { - logger.Log.Errorf(uiUtils.Error, fmt.Sprintf("Error getting pods by regex: %v", errormessage.FormatError(err))) - return - } - - if len(state.currentlyTappedPods) == 0 { - var suggestionStr string - if !shared.Contains(targetNamespaces, kubernetes.K8sAllNamespaces) { - suggestionStr = ". Select a different namespace with -n or tap all namespaces with -A" - } - logger.Log.Warningf(uiUtils.Warning, fmt.Sprintf("Did not find any pods matching the regex argument%s", suggestionStr)) - } - if config.Config.Tap.DryRun { return } @@ -146,14 +130,88 @@ func RunMizuTap() { } defer finishMizuExecution(kubernetesProvider) + if tapManagerErr := startTapManager(ctx, cancel, kubernetesProvider, targetNamespaces, *mizuApiFilteringOptions); tapManagerErr != nil { + logger.Log.Errorf(uiUtils.Error, getErrorDisplayTextForK8sTapManagerError(*tapManagerErr)) + cancel() + } + go goUtils.HandleExcWrapper(watchApiServerPod, ctx, kubernetesProvider, cancel, mizuApiFilteringOptions) go goUtils.HandleExcWrapper(watchTapperPod, ctx, kubernetesProvider, cancel) - go goUtils.HandleExcWrapper(watchPodsForTapping, ctx, kubernetesProvider, targetNamespaces, cancel, mizuApiFilteringOptions) // block until exit signal or error waitForFinish(ctx, cancel) } +func startTapManager(ctx context.Context, cancel context.CancelFunc, provider *kubernetes.Provider, targetNamespaces []string, mizuApiFilteringOptions api.TrafficFilteringOptions) *kubernetes.K8sTapManagerError { + manager, err := kubernetes.CreateAndStartK8sTapManager(ctx, provider, kubernetes.TapManagerConfig{ + TargetNamespaces: targetNamespaces, + PodFilterRegex: *config.Config.Tap.PodRegex(), + MizuResourcesNamespace: config.Config.MizuResourcesNamespace, + AgentImage: config.Config.AgentImage, + TapperResources: config.Config.Tap.TapperResources, + ImagePullPolicy: config.Config.ImagePullPolicy(), + DumpLogs: config.Config.DumpLogs, + IgnoredUserAgents: config.Config.Tap.IgnoredUserAgents, + MizuApiFilteringOptions: mizuApiFilteringOptions, + MizuServiceAccountExists: state.mizuServiceAccountExists, + }) + + if err != nil { + return err + } + + if len(manager.CurrentlyTappedPods) == 0 { + var suggestionStr string + if !shared.Contains(targetNamespaces, kubernetes.K8sAllNamespaces) { + suggestionStr = ". Select a different namespace with -n or tap all namespaces with -A" + } + logger.Log.Warningf(uiUtils.Warning, fmt.Sprintf("Did not find any pods matching the regex argument%s", suggestionStr)) + } + + go func() { + for { + select { + case managerErr := <-manager.ErrorOut: + logger.Log.Errorf(uiUtils.Error, getErrorDisplayTextForK8sTapManagerError(managerErr)) + cancel() + case tappedPodChanges := <- manager.TapPodChangesOut: + if err := apiserver.Provider.ReportTappedPods(manager.CurrentlyTappedPods); err != nil { + logger.Log.Debugf("[Error] failed update tapped pods %v", err) + } + displayTapPodChangesEvent(tappedPodChanges) + case <- ctx.Done(): + return + } + } + }() + + state.tapManager = manager + + return nil +} + +func displayTapPodChangesEvent(event kubernetes.TappedPodChangeEvent) { + for _, addedPod := range event.Added { + logger.Log.Infof(uiUtils.Green, fmt.Sprintf("+%s", addedPod.Name)) + } + for _, removedPod := range event.Removed { + logger.Log.Infof(uiUtils.Red, fmt.Sprintf("-%s", removedPod.Name)) + } +} + +func getErrorDisplayTextForK8sTapManagerError(err kubernetes.K8sTapManagerError) string { + switch err.TapManagerReason { + case kubernetes.TapManagerPodListError: + return fmt.Sprintf("Failed to update currently tapped pods: %v", err.OriginalError) + case kubernetes.TapManagerPodWatchError: + return fmt.Sprintf("Error occured in k8s pod watch: %v", err.OriginalError) + case kubernetes.TapManagerTapperUpdateError: + return fmt.Sprintf("Error updating tappers: %v", err.OriginalError) + default: + return fmt.Sprintf("Unknown error occured in k8s tap manager: %v", err.OriginalError) + } +} + func readValidationRules(file string) (string, error) { rules, err := shared.DecodeEnforcePolicy(file) if err != nil { @@ -266,48 +324,6 @@ func getSyncEntriesConfig() *shared.SyncEntriesConfig { } } -func updateMizuTappers(ctx context.Context, kubernetesProvider *kubernetes.Provider, mizuApiFilteringOptions *api.TrafficFilteringOptions) error { - nodeToTappedPodIPMap := kubernetes.GetNodeHostToTappedPodIpsMap(state.currentlyTappedPods) - - if len(nodeToTappedPodIPMap) > 0 { - var serviceAccountName string - if state.mizuServiceAccountExists { - serviceAccountName = kubernetes.ServiceAccountName - } else { - serviceAccountName = "" - } - - serializedMizuApiFilteringOptions, err := json.Marshal(mizuApiFilteringOptions) - if err != nil { - return err - } - - if err := kubernetesProvider.ApplyMizuTapperDaemonSet( - ctx, - config.Config.MizuResourcesNamespace, - kubernetes.TapperDaemonSetName, - config.Config.AgentImage, - kubernetes.TapperPodName, - fmt.Sprintf("%s.%s.svc.cluster.local", state.apiServerService.Name, state.apiServerService.Namespace), - nodeToTappedPodIPMap, - serviceAccountName, - config.Config.Tap.TapperResources, - config.Config.ImagePullPolicy(), - string(serializedMizuApiFilteringOptions), - config.Config.DumpLogs, - ); err != nil { - return err - } - logger.Log.Debugf("Successfully created %v tappers", len(nodeToTappedPodIPMap)) - } else { - if err := kubernetesProvider.RemoveDaemonSet(ctx, config.Config.MizuResourcesNamespace, kubernetes.TapperDaemonSetName); err != nil { - return err - } - } - - return nil -} - func finishMizuExecution(kubernetesProvider *kubernetes.Provider) { telemetry.ReportAPICalls() removalCtx, cancel := context.WithTimeout(context.Background(), cleanupTimeout) @@ -434,144 +450,6 @@ func waitUntilNamespaceDeleted(ctx context.Context, cancel context.CancelFunc, k } } -func watchPodsForTapping(ctx context.Context, kubernetesProvider *kubernetes.Provider, targetNamespaces []string, cancel context.CancelFunc, mizuApiFilteringOptions *api.TrafficFilteringOptions) { - added, modified, removed, errorChan := kubernetes.FilteredWatch(ctx, kubernetesProvider, targetNamespaces, config.Config.Tap.PodRegex()) - - restartTappers := func() { - err, changeFound := updateCurrentlyTappedPods(kubernetesProvider, ctx, targetNamespaces) - if err != nil { - logger.Log.Errorf(uiUtils.Error, fmt.Sprintf("Failed to update currently tapped pods: %v", err)) - cancel() - } - - if !changeFound { - logger.Log.Debugf("Nothing changed update tappers not needed") - return - } - - if err := apiserver.Provider.ReportTappedPods(state.currentlyTappedPods); err != nil { - logger.Log.Debugf("[Error] failed update tapped pods %v", err) - } - - if err := updateMizuTappers(ctx, kubernetesProvider, mizuApiFilteringOptions); err != nil { - logger.Log.Errorf(uiUtils.Error, fmt.Sprintf("Error updating tappers: %v", errormessage.FormatError(err))) - cancel() - } - } - restartTappersDebouncer := debounce.NewDebouncer(updateTappersDelay, restartTappers) - - for { - select { - case pod, ok := <-added: - if !ok { - added = nil - continue - } - - logger.Log.Debugf("Added matching pod %s, ns: %s", pod.Name, pod.Namespace) - restartTappersDebouncer.SetOn() - case pod, ok := <-removed: - if !ok { - removed = nil - continue - } - - logger.Log.Debugf("Removed matching pod %s, ns: %s", pod.Name, pod.Namespace) - restartTappersDebouncer.SetOn() - case pod, ok := <-modified: - if !ok { - modified = nil - continue - } - - logger.Log.Debugf("Modified matching pod %s, ns: %s, phase: %s, ip: %s", pod.Name, pod.Namespace, pod.Status.Phase, pod.Status.PodIP) - // Act only if the modified pod has already obtained an IP address. - // After filtering for IPs, on a normal pod restart this includes the following events: - // - Pod deletion - // - Pod reaches start state - // - Pod reaches ready state - // Ready/unready transitions might also trigger this event. - if pod.Status.PodIP != "" { - restartTappersDebouncer.SetOn() - } - case err, ok := <-errorChan: - if !ok { - errorChan = nil - continue - } - - logger.Log.Debugf("Watching pods loop, got error %v, stopping `restart tappers debouncer`", err) - restartTappersDebouncer.Cancel() - // TODO: Does this also perform cleanup? - cancel() - - case <-ctx.Done(): - logger.Log.Debugf("Watching pods loop, context done, stopping `restart tappers debouncer`") - restartTappersDebouncer.Cancel() - return - } - } -} - -func updateCurrentlyTappedPods(kubernetesProvider *kubernetes.Provider, ctx context.Context, targetNamespaces []string) (error, bool) { - changeFound := false - if matchingPods, err := kubernetesProvider.ListAllRunningPodsMatchingRegex(ctx, config.Config.Tap.PodRegex(), targetNamespaces); err != nil { - return err, false - } else { - podsToTap := excludeMizuPods(matchingPods) - addedPods, removedPods := getPodArrayDiff(state.currentlyTappedPods, podsToTap) - for _, addedPod := range addedPods { - changeFound = true - logger.Log.Infof(uiUtils.Green, fmt.Sprintf("+%s", addedPod.Name)) - } - for _, removedPod := range removedPods { - changeFound = true - logger.Log.Infof(uiUtils.Red, fmt.Sprintf("-%s", removedPod.Name)) - } - state.currentlyTappedPods = podsToTap - } - - return nil, changeFound -} - -func excludeMizuPods(pods []core.Pod) []core.Pod { - mizuPrefixRegex := regexp.MustCompile("^" + kubernetes.MizuResourcesPrefix) - - nonMizuPods := make([]core.Pod, 0) - for _, pod := range pods { - if !mizuPrefixRegex.MatchString(pod.Name) { - nonMizuPods = append(nonMizuPods, pod) - } - } - - return nonMizuPods -} - -func getPodArrayDiff(oldPods []core.Pod, newPods []core.Pod) (added []core.Pod, removed []core.Pod) { - added = getMissingPods(newPods, oldPods) - removed = getMissingPods(oldPods, newPods) - - return added, removed -} - -//returns pods present in pods1 array and missing in pods2 array -func getMissingPods(pods1 []core.Pod, pods2 []core.Pod) []core.Pod { - missingPods := make([]core.Pod, 0) - for _, pod1 := range pods1 { - var found = false - for _, pod2 := range pods2 { - if pod1.UID == pod2.UID { - found = true - break - } - } - if !found { - missingPods = append(missingPods, pod1) - } - } - return missingPods -} - func watchApiServerPod(ctx context.Context, kubernetesProvider *kubernetes.Provider, cancel context.CancelFunc, mizuApiFilteringOptions *api.TrafficFilteringOptions) { podExactRegex := regexp.MustCompile(fmt.Sprintf("^%s$", kubernetes.ApiServerPodName)) added, modified, removed, errorChan := kubernetes.FilteredWatch(ctx, kubernetesProvider, []string{config.Config.MizuResourcesNamespace}, podExactRegex) @@ -629,14 +507,10 @@ func watchApiServerPod(ctx context.Context, kubernetesProvider *kubernetes.Provi cancel() break } - if err := updateMizuTappers(ctx, kubernetesProvider, mizuApiFilteringOptions); err != nil { - logger.Log.Errorf(uiUtils.Error, fmt.Sprintf("Error updating tappers: %v", errormessage.FormatError(err))) - cancel() - } logger.Log.Infof("Mizu is available at %s\n", url) uiUtils.OpenBrowser(url) - if err := apiserver.Provider.ReportTappedPods(state.currentlyTappedPods); err != nil { + if err := apiserver.Provider.ReportTappedPods(state.tapManager.CurrentlyTappedPods); err != nil { logger.Log.Debugf("[Error] failed update tapped pods %v", err) } } @@ -646,7 +520,7 @@ func watchApiServerPod(ctx context.Context, kubernetesProvider *kubernetes.Provi continue } - logger.Log.Debugf("[ERROR] Agent creation, watching %v namespace, error: %v", config.Config.MizuResourcesNamespace, err) + logger.Log.Errorf("[ERROR] Agent creation, watching %v namespace, error: %v", config.Config.MizuResourcesNamespace, err) cancel() case <-timeAfter: @@ -717,7 +591,7 @@ func watchTapperPod(ctx context.Context, kubernetesProvider *kubernetes.Provider continue } - logger.Log.Debugf("[Error] Error in mizu tapper watch, err: %v", err) + logger.Log.Errorf("[Error] Error in mizu tapper watch, err: %v", err) cancel() case <-ctx.Done(): diff --git a/cli/mizu.go b/cli/mizu.go index 05d692c66..7c22213e7 100644 --- a/cli/mizu.go +++ b/cli/mizu.go @@ -2,7 +2,7 @@ package main import ( "github.com/up9inc/mizu/cli/cmd" - "github.com/up9inc/mizu/cli/mizu/goUtils" + "github.com/up9inc/mizu/shared/goUtils" ) func main() { diff --git a/shared/go.mod b/shared/go.mod index 8b0da2f05..539d4f50a 100644 --- a/shared/go.mod +++ b/shared/go.mod @@ -6,9 +6,12 @@ require ( github.com/docker/go-units v0.4.0 github.com/golang-jwt/jwt/v4 v4.1.0 github.com/op/go-logging v0.0.0-20160315200505-970db520ece7 + github.com/up9inc/mizu/tap/api v0.0.0 gopkg.in/yaml.v3 v3.0.0-20210107192922-496545a6307b k8s.io/api v0.21.0 k8s.io/apimachinery v0.21.0 k8s.io/client-go v0.21.0 k8s.io/kubectl v0.21.0 ) + +replace github.com/up9inc/mizu/tap/api v0.0.0 => ../tap/api diff --git a/shared/go.sum b/shared/go.sum index 5ccae47fe..20c15962c 100644 --- a/shared/go.sum +++ b/shared/go.sum @@ -211,6 +211,7 @@ github.com/google/go-cmp v0.5.2/go.mod h1:v8dTdLbMG2kIc/vJvl+f65V22dbkXbowE6jgT/ github.com/google/gofuzz v1.0.0/go.mod h1:dBl0BpW6vV/+mYPU4Po3pmUjxk6FQPldtuIdl/M65Eg= github.com/google/gofuzz v1.1.0 h1:Hsa8mG0dQ46ij8Sl2AYJDUv1oA9/d6Vk+3LG99Oe02g= github.com/google/gofuzz v1.1.0/go.mod h1:dBl0BpW6vV/+mYPU4Po3pmUjxk6FQPldtuIdl/M65Eg= +github.com/google/martian v2.1.0+incompatible h1:/CP5g8u/VJHijgedC/Legn3BAbAaWPgecwXBIDzw5no= github.com/google/martian v2.1.0+incompatible/go.mod h1:9I4somxYTbIHy5NJKHRl3wXiIaQGbYVAs8BPL6v8lEs= github.com/google/pprof v0.0.0-20181206194817-3ea8567a2e57/go.mod h1:zfwlbNMJ+OItoe0UupaVj+oy1omPYYDuagoSzA8v9mc= github.com/google/pprof v0.0.0-20190515194954-54271f7e092f/go.mod h1:zfwlbNMJ+OItoe0UupaVj+oy1omPYYDuagoSzA8v9mc= diff --git a/cli/mizu/goUtils/funcWrappers.go b/shared/goUtils/funcWrappers.go similarity index 100% rename from cli/mizu/goUtils/funcWrappers.go rename to shared/goUtils/funcWrappers.go diff --git a/shared/kubernetes/k8sTapManager.go b/shared/kubernetes/k8sTapManager.go index 20315e7dd..22e8d2e30 100644 --- a/shared/kubernetes/k8sTapManager.go +++ b/shared/kubernetes/k8sTapManager.go @@ -3,96 +3,209 @@ package kubernetes import ( "context" "fmt" - "github.com/up9inc/mizu/cli/uiUtils" "github.com/up9inc/mizu/shared" + "github.com/up9inc/mizu/shared/debounce" + "github.com/up9inc/mizu/shared/goUtils" "github.com/up9inc/mizu/shared/logger" + "github.com/up9inc/mizu/tap/api" core "k8s.io/api/core/v1" "regexp" + "time" ) -type K8sTapManager struct { - state TapState - config TapManagerConfig - podFilterRegex regexp.Regexp +const updateTappersDelay = 5 * time.Second + +type TappedPodChangeEvent struct { + Added []core.Pod + Removed []core.Pod } -type TapState struct { - apiServerService *core.Service - currentlyTappedPods []core.Pod - mizuServiceAccountExists bool +type K8sTapManager struct { + context context.Context + CurrentlyTappedPods []core.Pod + config TapManagerConfig + kubernetesProvider *Provider + TapPodChangesOut chan TappedPodChangeEvent + ErrorOut chan K8sTapManagerError } type TapManagerConfig struct { - MizuResourcesNamespace string - AgentImage string - TapperResources shared.Resources - ImagePullPolicy core.PullPolicy - DumpLogs bool - IgnoredUserAgents []string + TargetNamespaces []string + PodFilterRegex regexp.Regexp + MizuResourcesNamespace string + AgentImage string + TapperResources shared.Resources + ImagePullPolicy core.PullPolicy + DumpLogs bool + IgnoredUserAgents []string + MizuApiFilteringOptions api.TrafficFilteringOptions + MizuServiceAccountExists bool } -func CreateK8sTapManager(config TapManagerConfig, podFilterRegex regexp.Regexp) *K8sTapManager { - return &K8sTapManager{ - state: TapState{}, - config: config, - podFilterRegex: podFilterRegex, +func CreateAndStartK8sTapManager(ctx context.Context, kubernetesProvider *Provider, config TapManagerConfig) (*K8sTapManager, *K8sTapManagerError) { + manager := &K8sTapManager{ + context: ctx, + CurrentlyTappedPods: make([]core.Pod, 0), + config: config, + kubernetesProvider: kubernetesProvider, + TapPodChangesOut: make(chan TappedPodChangeEvent, 100), + ErrorOut: make(chan K8sTapManagerError, 100), + } + + err, _ := manager.updateCurrentlyTappedPods() + if err != nil { + return nil, &K8sTapManagerError{ + OriginalError: err, + TapManagerReason: TapManagerPodListError, + } + } + + err = manager.updateMizuTappers() + if err != nil { + return nil, &K8sTapManagerError{ + OriginalError: err, + TapManagerReason: TapManagerTapperUpdateError, + } + } + + go goUtils.HandleExcWrapper(manager.watchPodsForTapping) + return manager, nil +} + + +func (tapManager *K8sTapManager) watchPodsForTapping() { + added, modified, removed, errorChan := FilteredWatch(tapManager.context, tapManager.kubernetesProvider, tapManager.config.TargetNamespaces, &tapManager.config.PodFilterRegex) + + restartTappers := func() { + err, changeFound := tapManager.updateCurrentlyTappedPods() + if err != nil { + tapManager.ErrorOut <- K8sTapManagerError{ + OriginalError: err, + TapManagerReason: TapManagerPodListError, + } + } + + if !changeFound { + logger.Log.Debugf("Nothing changed update tappers not needed") + return + } + + if err := tapManager.updateMizuTappers(); err != nil { + tapManager.ErrorOut <- K8sTapManagerError{ + OriginalError: err, + TapManagerReason: TapManagerTapperUpdateError, + } + } + } + restartTappersDebouncer := debounce.NewDebouncer(updateTappersDelay, restartTappers) + + for { + select { + case pod, ok := <-added: + if !ok { + added = nil + continue + } + + logger.Log.Debugf("Added matching pod %s, ns: %s", pod.Name, pod.Namespace) + restartTappersDebouncer.SetOn() + case pod, ok := <-removed: + if !ok { + removed = nil + continue + } + + logger.Log.Debugf("Removed matching pod %s, ns: %s", pod.Name, pod.Namespace) + restartTappersDebouncer.SetOn() + case pod, ok := <-modified: + if !ok { + modified = nil + continue + } + + logger.Log.Debugf("Modified matching pod %s, ns: %s, phase: %s, ip: %s", pod.Name, pod.Namespace, pod.Status.Phase, pod.Status.PodIP) + // Act only if the modified pod has already obtained an IP address. + // After filtering for IPs, on a normal pod restart this includes the following events: + // - Pod deletion + // - Pod reaches start state + // - Pod reaches ready state + // Ready/unready transitions might also trigger this event. + if pod.Status.PodIP != "" { + restartTappersDebouncer.SetOn() + } + case err, ok := <-errorChan: + if !ok { + errorChan = nil + continue + } + + logger.Log.Debugf("Watching pods loop, got error %v, stopping `restart tappers debouncer`", err) + restartTappersDebouncer.Cancel() + tapManager.ErrorOut <- K8sTapManagerError{ + OriginalError: err, + TapManagerReason: TapManagerPodWatchError, + } + // TODO: Does this also perform cleanup? + + case <-tapManager.context.Done(): + logger.Log.Debugf("Watching pods loop, context done, stopping `restart tappers debouncer`") + restartTappersDebouncer.Cancel() + return + } } } -func updateCurrentlyTappedPods(kubernetesProvider *Provider, ctx context.Context, targetNamespaces []string) (error, bool) { - changeFound := false - if matchingPods, err := kubernetesProvider.ListAllRunningPodsMatchingRegex(ctx, config.Config.Tap.PodRegex(), targetNamespaces); err != nil { +func (tapManager *K8sTapManager) updateCurrentlyTappedPods() (err error, changesFound bool) { + if matchingPods, err := tapManager.kubernetesProvider.ListAllRunningPodsMatchingRegex(tapManager.context, &tapManager.config.PodFilterRegex, tapManager.config.TargetNamespaces); err != nil { return err, false } else { podsToTap := excludeMizuPods(matchingPods) - addedPods, removedPods := getPodArrayDiff(state.currentlyTappedPods, podsToTap) - for _, addedPod := range addedPods { - changeFound = true - logger.Log.Infof(uiUtils.Green, fmt.Sprintf("+%s", addedPod.Name)) + addedPods, removedPods := getPodArrayDiff(tapManager.CurrentlyTappedPods, podsToTap) + if len(addedPods) > 0 || len(removedPods) > 0 { + tapManager.CurrentlyTappedPods = podsToTap + tapManager.TapPodChangesOut <- TappedPodChangeEvent{ + Added: addedPods, + Removed: removedPods, + } + return nil, true } - for _, removedPod := range removedPods { - changeFound = true - logger.Log.Infof(uiUtils.Red, fmt.Sprintf("-%s", removedPod.Name)) - } - state.currentlyTappedPods = podsToTap + return nil, false } - - return nil, changeFound } -func (tapManager *K8sTapManager) updateMizuTappers(ctx context.Context, kubernetesProvider *Provider, mizuApiFilteringOptions *interface{}) error { - nodeToTappedPodIPMap := GetNodeHostToTappedPodIpsMap(tapManager.state.currentlyTappedPods) +func (tapManager *K8sTapManager) updateMizuTappers() error { + nodeToTappedPodIPMap := GetNodeHostToTappedPodIpsMap(tapManager.CurrentlyTappedPods) if len(nodeToTappedPodIPMap) > 0 { var serviceAccountName string - if tapManager.state.mizuServiceAccountExists { + if tapManager.config.MizuServiceAccountExists { serviceAccountName = ServiceAccountName } else { serviceAccountName = "" } - if err := kubernetesProvider.ApplyMizuTapperDaemonSet( - ctx, + if err := tapManager.kubernetesProvider.ApplyMizuTapperDaemonSet( + tapManager.context, tapManager.config.MizuResourcesNamespace, TapperDaemonSetName, tapManager.config.AgentImage, TapperPodName, - fmt.Sprintf("%s.%s.svc.cluster.local", tapManager.state.apiServerService.Name, tapManager.state.apiServerService.Namespace), + fmt.Sprintf("%s.%s.svc.cluster.local", ApiServerPodName, tapManager.config.MizuResourcesNamespace), nodeToTappedPodIPMap, serviceAccountName, tapManager.config.TapperResources, tapManager.config.ImagePullPolicy, - mizuApiFilteringOptions, + tapManager.config.MizuApiFilteringOptions, tapManager.config.DumpLogs, ); err != nil { return err } logger.Log.Debugf("Successfully created %v tappers", len(nodeToTappedPodIPMap)) } else { - if err := kubernetesProvider.RemoveDaemonSet(ctx, tapManager.config.MizuResourcesNamespace, TapperDaemonSetName); err != nil { + if err := tapManager.kubernetesProvider.RemoveDaemonSet(tapManager.context, tapManager.config.MizuResourcesNamespace, TapperDaemonSetName); err != nil { return err } } return nil -} \ No newline at end of file +} diff --git a/shared/kubernetes/k8sTapManagerErrors.go b/shared/kubernetes/k8sTapManagerErrors.go new file mode 100644 index 000000000..47eec3efe --- /dev/null +++ b/shared/kubernetes/k8sTapManagerErrors.go @@ -0,0 +1,19 @@ +package kubernetes + +type K8sTapManagerErrorReason string + +const ( + TapManagerTapperUpdateError K8sTapManagerErrorReason = "TAPPER_UPDATE_ERROR" + TapManagerPodWatchError K8sTapManagerErrorReason = "POD_WATCH_ERROR" + TapManagerPodListError K8sTapManagerErrorReason = "POD_LIST_ERROR" +) + +type K8sTapManagerError struct { + OriginalError error + TapManagerReason K8sTapManagerErrorReason +} + +// K8sTapManagerError implements the Error interface. +func (e *K8sTapManagerError) Error() string { + return e.OriginalError.Error() +} diff --git a/shared/kubernetes/provider.go b/shared/kubernetes/provider.go index e0c111eaf..87dc52617 100644 --- a/shared/kubernetes/provider.go +++ b/shared/kubernetes/provider.go @@ -9,6 +9,7 @@ import ( "github.com/up9inc/mizu/shared" "github.com/up9inc/mizu/shared/logger" "github.com/up9inc/mizu/shared/semver" + "github.com/up9inc/mizu/tap/api" "io" core "k8s.io/api/core/v1" rbac "k8s.io/api/rbac/v1" @@ -510,7 +511,7 @@ func (provider *Provider) CreateConfigMap(ctx context.Context, namespace string, return nil } -func (provider *Provider) ApplyMizuTapperDaemonSet(ctx context.Context, namespace string, daemonSetName string, podImage string, tapperPodName string, apiServerPodIp string, nodeToTappedPodIPMap map[string][]string, serviceAccountName string, resources shared.Resources, imagePullPolicy core.PullPolicy, mizuApiFilteringOptions *interface{}, dumpLogs bool) error { +func (provider *Provider) ApplyMizuTapperDaemonSet(ctx context.Context, namespace string, daemonSetName string, podImage string, tapperPodName string, apiServerPodIp string, nodeToTappedPodIPMap map[string][]string, serviceAccountName string, resources shared.Resources, imagePullPolicy core.PullPolicy, mizuApiFilteringOptions api.TrafficFilteringOptions, dumpLogs bool) error { logger.Log.Debugf("Applying %d tapper daemon sets, ns: %s, daemonSetName: %s, podImage: %s, tapperPodName: %s", len(nodeToTappedPodIPMap), namespace, daemonSetName, podImage, tapperPodName) if len(nodeToTappedPodIPMap) == 0 { diff --git a/shared/kubernetes/utils.go b/shared/kubernetes/utils.go index a50adc3fd..4c7697d15 100644 --- a/shared/kubernetes/utils.go +++ b/shared/kubernetes/utils.go @@ -1,6 +1,9 @@ package kubernetes -import core "k8s.io/api/core/v1" +import ( + core "k8s.io/api/core/v1" + "regexp" +) func GetNodeHostToTappedPodIpsMap(tappedPods []core.Pod) map[string][]string { nodeToTappedPodIPMap := make(map[string][]string, 0) @@ -13,4 +16,43 @@ func GetNodeHostToTappedPodIpsMap(tappedPods []core.Pod) map[string][]string { } } return nodeToTappedPodIPMap +} + + +func excludeMizuPods(pods []core.Pod) []core.Pod { + mizuPrefixRegex := regexp.MustCompile("^" + MizuResourcesPrefix) + + nonMizuPods := make([]core.Pod, 0) + for _, pod := range pods { + if !mizuPrefixRegex.MatchString(pod.Name) { + nonMizuPods = append(nonMizuPods, pod) + } + } + + return nonMizuPods +} + +func getPodArrayDiff(oldPods []core.Pod, newPods []core.Pod) (added []core.Pod, removed []core.Pod) { + added = getMissingPods(newPods, oldPods) + removed = getMissingPods(oldPods, newPods) + + return added, removed +} + +//returns pods present in pods1 array and missing in pods2 array +func getMissingPods(pods1 []core.Pod, pods2 []core.Pod) []core.Pod { + missingPods := make([]core.Pod, 0) + for _, pod1 := range pods1 { + var found = false + for _, pod2 := range pods2 { + if pod1.UID == pod2.UID { + found = true + break + } + } + if !found { + missingPods = append(missingPods, pod1) + } + } + return missingPods } \ No newline at end of file