diff --git a/cli/cmd/tapRunner.go b/cli/cmd/tapRunner.go index 29f6e1f47..baa2c4dc7 100644 --- a/cli/cmd/tapRunner.go +++ b/cli/cmd/tapRunner.go @@ -33,13 +33,11 @@ import ( "github.com/up9inc/mizu/tap/api" ) -const ( - cleanupTimeout = time.Minute -) +const cleanupTimeout = time.Minute type tapState struct { apiServerService *core.Service - tapManager *kubernetes.K8sTapManager + tapperSyncer *kubernetes.MizuTapperSyncer mizuServiceAccountExists bool } @@ -149,7 +147,7 @@ func RunMizuTap() { } func startTapManager(ctx context.Context, cancel context.CancelFunc, provider *kubernetes.Provider, targetNamespaces []string, mizuApiFilteringOptions api.TrafficFilteringOptions) error { - manager, err := kubernetes.CreateAndStartK8sTapManager(ctx, provider, kubernetes.TapManagerConfig{ + tapperSyncer, err := kubernetes.CreateAndStartMizuTapperSyncer(ctx, provider, kubernetes.TapperSyncerConfig{ TargetNamespaces: targetNamespaces, PodFilterRegex: *config.Config.Tap.PodRegex(), MizuResourcesNamespace: config.Config.MizuResourcesNamespace, @@ -166,7 +164,7 @@ func startTapManager(ctx context.Context, cancel context.CancelFunc, provider *k return err } - if len(manager.CurrentlyTappedPods) == 0 { + if len(tapperSyncer.CurrentlyTappedPods) == 0 { var suggestionStr string if !shared.Contains(targetNamespaces, kubernetes.K8sAllNamespaces) { suggestionStr = ". Select a different namespace with -n or tap all namespaces with -A" @@ -177,21 +175,21 @@ func startTapManager(ctx context.Context, cancel context.CancelFunc, provider *k go func() { for { select { - case managerErr := <-manager.ErrorOut: + case managerErr := <-tapperSyncer.ErrorOut: logger.Log.Errorf(uiUtils.Error, getErrorDisplayTextForK8sTapManagerError(managerErr)) cancel() - case tappedPodChanges := <- manager.TapPodChangesOut: - if err := apiserver.Provider.ReportTappedPods(manager.CurrentlyTappedPods); err != nil { + case tappedPodChanges := <-tapperSyncer.TapPodChangesOut: + if err := apiserver.Provider.ReportTappedPods(tapperSyncer.CurrentlyTappedPods); err != nil { logger.Log.Debugf("[Error] failed update tapped pods %v", err) } displayTapPodChangesEvent(tappedPodChanges) - case <- ctx.Done(): + case <-ctx.Done(): return } } }() - state.tapManager = manager + state.tapperSyncer = tapperSyncer return nil } @@ -514,7 +512,7 @@ func watchApiServerPod(ctx context.Context, kubernetesProvider *kubernetes.Provi break } - if err := state.tapManager.BeginUpdatingTappers(); err != nil { + if err := state.tapperSyncer.BeginUpdatingTappers(); err != nil { logger.Log.Errorf(uiUtils.Error, fmt.Sprintf("Error updating tappers: %v", err)) cancel() break @@ -522,7 +520,7 @@ func watchApiServerPod(ctx context.Context, kubernetesProvider *kubernetes.Provi logger.Log.Infof("Mizu is available at %s\n", url) uiUtils.OpenBrowser(url) - if err := apiserver.Provider.ReportTappedPods(state.tapManager.CurrentlyTappedPods); err != nil { + if err := apiserver.Provider.ReportTappedPods(state.tapperSyncer.CurrentlyTappedPods); err != nil { logger.Log.Debugf("[Error] failed update tapped pods %v", err) } } diff --git a/cli/config/config.go b/cli/config/config.go index a664d66b4..753b7f429 100644 --- a/cli/config/config.go +++ b/cli/config/config.go @@ -84,6 +84,7 @@ func WriteConfig(config *ConfigStruct) error { } type updateConfigStruct func(*ConfigStruct) + func UpdateConfig(updateConfigStruct updateConfigStruct) error { configFile, err := GetConfigWithDefaults() if err != nil { diff --git a/cli/config/configStructs/logsConfig.go b/cli/config/configStructs/logsConfig.go index 90d41888f..14c4c12b4 100644 --- a/cli/config/configStructs/logsConfig.go +++ b/cli/config/configStructs/logsConfig.go @@ -32,4 +32,4 @@ func (config *LogsConfig) FilePath() string { } return config.FileStr -} \ No newline at end of file +} diff --git a/shared/kubernetes/consts.go b/shared/kubernetes/consts.go index b52e4539c..d95ec3b0f 100644 --- a/shared/kubernetes/consts.go +++ b/shared/kubernetes/consts.go @@ -14,4 +14,3 @@ const ( ConfigMapName = MizuResourcesPrefix + "config" MinKubernetesServerVersion = "1.16.0" ) - diff --git a/shared/kubernetes/errors.go b/shared/kubernetes/errors.go index cfd26e2ef..5b65e2d17 100644 --- a/shared/kubernetes/errors.go +++ b/shared/kubernetes/errors.go @@ -4,12 +4,12 @@ type K8sTapManagerErrorReason string const ( TapManagerTapperUpdateError K8sTapManagerErrorReason = "TAPPER_UPDATE_ERROR" - TapManagerPodWatchError K8sTapManagerErrorReason = "POD_WATCH_ERROR" - TapManagerPodListError K8sTapManagerErrorReason = "POD_LIST_ERROR" + TapManagerPodWatchError K8sTapManagerErrorReason = "POD_WATCH_ERROR" + TapManagerPodListError K8sTapManagerErrorReason = "POD_LIST_ERROR" ) type K8sTapManagerError struct { - OriginalError error + OriginalError error TapManagerReason K8sTapManagerErrorReason } @@ -18,7 +18,7 @@ func (e *K8sTapManagerError) Error() string { return e.OriginalError.Error() } -type ClusterBehindProxyError struct {} +type ClusterBehindProxyError struct{} // ClusterBehindProxyError implements the Error interface. func (e *ClusterBehindProxyError) Error() string { diff --git a/shared/kubernetes/k8sTapManager.go b/shared/kubernetes/k8sTapManager.go index 7a83b668f..06dab49fd 100644 --- a/shared/kubernetes/k8sTapManager.go +++ b/shared/kubernetes/k8sTapManager.go @@ -20,17 +20,18 @@ type TappedPodChangeEvent struct { Removed []core.Pod } -type K8sTapManager struct { +// MizuTapperSyncer syncs tappers using a k8s pod watch +type MizuTapperSyncer struct { context context.Context CurrentlyTappedPods []core.Pod - config TapManagerConfig + config TapperSyncerConfig kubernetesProvider *Provider TapPodChangesOut chan TappedPodChangeEvent ErrorOut chan K8sTapManagerError - shouldUpdateTappers bool // used to prevent updating tapper daemonsets before api is available + shouldUpdateTappers bool // used to prevent updating tapper daemonsets before api is available but still get targeted pod change events } -type TapManagerConfig struct { +type TapperSyncerConfig struct { TargetNamespaces []string PodFilterRegex regexp.Regexp MizuResourcesNamespace string @@ -43,8 +44,8 @@ type TapManagerConfig struct { MizuServiceAccountExists bool } -func CreateAndStartK8sTapManager(ctx context.Context, kubernetesProvider *Provider, config TapManagerConfig, shouldUpdateTappers bool) (*K8sTapManager, error) { - manager := &K8sTapManager{ +func CreateAndStartMizuTapperSyncer(ctx context.Context, kubernetesProvider *Provider, config TapperSyncerConfig, shouldUpdateTappers bool) (*MizuTapperSyncer, error) { + manager := &MizuTapperSyncer{ context: ctx, CurrentlyTappedPods: make([]core.Pod, 0), config: config, @@ -69,21 +70,21 @@ func CreateAndStartK8sTapManager(ctx context.Context, kubernetesProvider *Provid } // BeginUpdatingTappers should only be called after mizu api server is available -func (tapManager *K8sTapManager) BeginUpdatingTappers() error { - tapManager.shouldUpdateTappers = true - if err := tapManager.updateMizuTappers(); err != nil { +func (tapperSyncer *MizuTapperSyncer) BeginUpdatingTappers() error { + tapperSyncer.shouldUpdateTappers = true + if err := tapperSyncer.updateMizuTappers(); err != nil { return err } return nil } -func (tapManager *K8sTapManager) watchPodsForTapping() { - added, modified, removed, errorChan := FilteredWatch(tapManager.context, tapManager.kubernetesProvider, tapManager.config.TargetNamespaces, &tapManager.config.PodFilterRegex) +func (tapperSyncer *MizuTapperSyncer) watchPodsForTapping() { + added, modified, removed, errorChan := FilteredWatch(tapperSyncer.context, tapperSyncer.kubernetesProvider, tapperSyncer.config.TargetNamespaces, &tapperSyncer.config.PodFilterRegex) restartTappers := func() { - err, changeFound := tapManager.updateCurrentlyTappedPods() + err, changeFound := tapperSyncer.updateCurrentlyTappedPods() if err != nil { - tapManager.ErrorOut <- K8sTapManagerError{ + tapperSyncer.ErrorOut <- K8sTapManagerError{ OriginalError: err, TapManagerReason: TapManagerPodListError, } @@ -93,9 +94,9 @@ func (tapManager *K8sTapManager) watchPodsForTapping() { logger.Log.Debugf("Nothing changed update tappers not needed") return } - if tapManager.shouldUpdateTappers { - if err := tapManager.updateMizuTappers(); err != nil { - tapManager.ErrorOut <- K8sTapManagerError{ + if tapperSyncer.shouldUpdateTappers { + if err := tapperSyncer.updateMizuTappers(); err != nil { + tapperSyncer.ErrorOut <- K8sTapManagerError{ OriginalError: err, TapManagerReason: TapManagerTapperUpdateError, } @@ -146,12 +147,12 @@ func (tapManager *K8sTapManager) watchPodsForTapping() { logger.Log.Debugf("Watching pods loop, got error %v, stopping `restart tappers debouncer`", err) restartTappersDebouncer.Cancel() - tapManager.ErrorOut <- K8sTapManagerError{ + tapperSyncer.ErrorOut <- K8sTapManagerError{ OriginalError: err, TapManagerReason: TapManagerPodWatchError, } - case <-tapManager.context.Done(): + case <-tapperSyncer.context.Done(): logger.Log.Debugf("Watching pods loop, context done, stopping `restart tappers debouncer`") restartTappersDebouncer.Cancel() return @@ -159,15 +160,15 @@ func (tapManager *K8sTapManager) watchPodsForTapping() { } } -func (tapManager *K8sTapManager) updateCurrentlyTappedPods() (err error, changesFound bool) { - if matchingPods, err := tapManager.kubernetesProvider.ListAllRunningPodsMatchingRegex(tapManager.context, &tapManager.config.PodFilterRegex, tapManager.config.TargetNamespaces); err != nil { +func (tapperSyncer *MizuTapperSyncer) updateCurrentlyTappedPods() (err error, changesFound bool) { + if matchingPods, err := tapperSyncer.kubernetesProvider.ListAllRunningPodsMatchingRegex(tapperSyncer.context, &tapperSyncer.config.PodFilterRegex, tapperSyncer.config.TargetNamespaces); err != nil { return err, false } else { podsToTap := excludeMizuPods(matchingPods) - addedPods, removedPods := getPodArrayDiff(tapManager.CurrentlyTappedPods, podsToTap) + addedPods, removedPods := getPodArrayDiff(tapperSyncer.CurrentlyTappedPods, podsToTap) if len(addedPods) > 0 || len(removedPods) > 0 { - tapManager.CurrentlyTappedPods = podsToTap - tapManager.TapPodChangesOut <- TappedPodChangeEvent{ + tapperSyncer.CurrentlyTappedPods = podsToTap + tapperSyncer.TapPodChangesOut <- TappedPodChangeEvent{ Added: addedPods, Removed: removedPods, } @@ -177,36 +178,36 @@ func (tapManager *K8sTapManager) updateCurrentlyTappedPods() (err error, changes } } -func (tapManager *K8sTapManager) updateMizuTappers() error { - nodeToTappedPodIPMap := GetNodeHostToTappedPodIpsMap(tapManager.CurrentlyTappedPods) +func (tapperSyncer *MizuTapperSyncer) updateMizuTappers() error { + nodeToTappedPodIPMap := GetNodeHostToTappedPodIpsMap(tapperSyncer.CurrentlyTappedPods) if len(nodeToTappedPodIPMap) > 0 { var serviceAccountName string - if tapManager.config.MizuServiceAccountExists { + if tapperSyncer.config.MizuServiceAccountExists { serviceAccountName = ServiceAccountName } else { serviceAccountName = "" } - if err := tapManager.kubernetesProvider.ApplyMizuTapperDaemonSet( - tapManager.context, - tapManager.config.MizuResourcesNamespace, + if err := tapperSyncer.kubernetesProvider.ApplyMizuTapperDaemonSet( + tapperSyncer.context, + tapperSyncer.config.MizuResourcesNamespace, TapperDaemonSetName, - tapManager.config.AgentImage, + tapperSyncer.config.AgentImage, TapperPodName, - fmt.Sprintf("%s.%s.svc.cluster.local", ApiServerPodName, tapManager.config.MizuResourcesNamespace), + fmt.Sprintf("%s.%s.svc.cluster.local", ApiServerPodName, tapperSyncer.config.MizuResourcesNamespace), nodeToTappedPodIPMap, serviceAccountName, - tapManager.config.TapperResources, - tapManager.config.ImagePullPolicy, - tapManager.config.MizuApiFilteringOptions, - tapManager.config.DumpLogs, + tapperSyncer.config.TapperResources, + tapperSyncer.config.ImagePullPolicy, + tapperSyncer.config.MizuApiFilteringOptions, + tapperSyncer.config.DumpLogs, ); err != nil { return err } logger.Log.Debugf("Successfully created %v tappers", len(nodeToTappedPodIPMap)) } else { - if err := tapManager.kubernetesProvider.RemoveDaemonSet(tapManager.context, tapManager.config.MizuResourcesNamespace, TapperDaemonSetName); err != nil { + if err := tapperSyncer.kubernetesProvider.RemoveDaemonSet(tapperSyncer.context, tapperSyncer.config.MizuResourcesNamespace, TapperDaemonSetName); err != nil { return err } } diff --git a/shared/kubernetes/provider.go b/shared/kubernetes/provider.go index 9ff029584..bccdcf1d7 100644 --- a/shared/kubernetes/provider.go +++ b/shared/kubernetes/provider.go @@ -609,7 +609,7 @@ func (provider *Provider) ApplyMizuTapperDaemonSet(ctx context.Context, namespac volumeName := ConfigMapName configMapVolume := applyconfcore.VolumeApplyConfiguration{ - Name: &volumeName, + Name: &volumeName, VolumeSourceApplyConfiguration: applyconfcore.VolumeSourceApplyConfiguration{ ConfigMap: &applyconfcore.ConfigMapVolumeSourceApplyConfiguration{ LocalObjectReferenceApplyConfiguration: applyconfcore.LocalObjectReferenceApplyConfiguration{ @@ -620,8 +620,8 @@ func (provider *Provider) ApplyMizuTapperDaemonSet(ctx context.Context, namespac } mountPath := shared.ConfigDirPath configMapVolumeMount := applyconfcore.VolumeMountApplyConfiguration{ - Name: &volumeName, - MountPath: &mountPath, + Name: &volumeName, + MountPath: &mountPath, } agentContainer.WithVolumeMounts(&configMapVolumeMount) @@ -773,7 +773,6 @@ func validateNotProxy(kubernetesConfig clientcmd.ClientConfig, restClientConfig return nil } - func validateKubernetesVersion(clientSet *kubernetes.Clientset) error { serverVersion, err := clientSet.ServerVersion() if err != nil { diff --git a/shared/kubernetes/utils.go b/shared/kubernetes/utils.go index 2eb970308..f77f27f27 100644 --- a/shared/kubernetes/utils.go +++ b/shared/kubernetes/utils.go @@ -18,7 +18,6 @@ func GetNodeHostToTappedPodIpsMap(tappedPods []core.Pod) map[string][]string { return nodeToTappedPodIPMap } - func excludeMizuPods(pods []core.Pod) []core.Pod { mizuPrefixRegex := regexp.MustCompile("^" + MizuResourcesPrefix) diff --git a/shared/kubernetes/watch.go b/shared/kubernetes/watch.go index 37e5bd576..58460f584 100644 --- a/shared/kubernetes/watch.go +++ b/shared/kubernetes/watch.go @@ -21,7 +21,6 @@ func FilteredWatch(ctx context.Context, kubernetesProvider *Provider, targetName removedChan := make(chan *corev1.Pod) errorChan := make(chan error) - var wg sync.WaitGroup for _, targetNamespace := range targetNamespaces { @@ -29,7 +28,7 @@ func FilteredWatch(ctx context.Context, kubernetesProvider *Provider, targetName go func(targetNamespace string) { defer wg.Done() - watchRestartDebouncer := debounce.NewDebouncer(1 * time.Minute, func() {}) + watchRestartDebouncer := debounce.NewDebouncer(1*time.Minute, func() {}) for { watcher := kubernetesProvider.GetPodWatcher(ctx, targetNamespace) @@ -37,7 +36,7 @@ func FilteredWatch(ctx context.Context, kubernetesProvider *Provider, targetName watcher.Stop() select { - case <- ctx.Done(): + case <-ctx.Done(): return default: break