diff --git a/cli/cmd/tapRunner.go b/cli/cmd/tapRunner.go index ff20e1287..bed740029 100644 --- a/cli/cmd/tapRunner.go +++ b/cli/cmd/tapRunner.go @@ -267,7 +267,7 @@ func getSyncEntriesConfig() *shared.SyncEntriesConfig { } func updateMizuTappers(ctx context.Context, kubernetesProvider *kubernetes.Provider, mizuApiFilteringOptions *api.TrafficFilteringOptions) error { - nodeToTappedPodIPMap := getNodeHostToTappedPodIpsMap(state.currentlyTappedPods) + nodeToTappedPodIPMap := kubernetes.GetNodeHostToTappedPodIpsMap(state.currentlyTappedPods) if len(nodeToTappedPodIPMap) > 0 { var serviceAccountName string @@ -742,19 +742,6 @@ func createRBACIfNecessary(ctx context.Context, kubernetesProvider *kubernetes.P return true, nil } -func getNodeHostToTappedPodIpsMap(tappedPods []core.Pod) map[string][]string { - nodeToTappedPodIPMap := make(map[string][]string, 0) - for _, pod := range tappedPods { - existingList := nodeToTappedPodIPMap[pod.Spec.NodeName] - if existingList == nil { - nodeToTappedPodIPMap[pod.Spec.NodeName] = []string{pod.Status.PodIP} - } else { - nodeToTappedPodIPMap[pod.Spec.NodeName] = append(nodeToTappedPodIPMap[pod.Spec.NodeName], pod.Status.PodIP) - } - } - return nodeToTappedPodIPMap -} - func getNamespaces(kubernetesProvider *kubernetes.Provider) []string { if config.Config.Tap.AllNamespaces { return []string{kubernetes.K8sAllNamespaces} diff --git a/shared/kubernetes/k8sTapManagement.go b/shared/kubernetes/k8sTapManagement.go new file mode 100644 index 000000000..d86f89e0a --- /dev/null +++ b/shared/kubernetes/k8sTapManagement.go @@ -0,0 +1,66 @@ +package kubernetes + +import ( + "context" + "fmt" + "github.com/up9inc/mizu/shared" + "github.com/up9inc/mizu/shared/logger" + core "k8s.io/api/core/v1" +) + +type K8sTapManager struct { + state TapState + Config TapManagerConfig +} + +type TapState struct { + apiServerService *core.Service + currentlyTappedPods []core.Pod + mizuServiceAccountExists bool +} + +type TapManagerConfig struct { + MizuResourcesNamespace string + AgentImage string + TapperResources shared.Resources + ImagePullPolicy core.PullPolicy + DumpLogs bool + +} + +func (tapManager *K8sTapManager) updateMizuTappers(ctx context.Context, kubernetesProvider *Provider, mizuApiFilteringOptions *interface{}) error { + nodeToTappedPodIPMap := GetNodeHostToTappedPodIpsMap(tapManager.state.currentlyTappedPods) + + if len(nodeToTappedPodIPMap) > 0 { + var serviceAccountName string + if tapManager.state.mizuServiceAccountExists { + serviceAccountName = ServiceAccountName + } else { + serviceAccountName = "" + } + + if err := kubernetesProvider.ApplyMizuTapperDaemonSet( + ctx, + tapManager.Config.MizuResourcesNamespace, + TapperDaemonSetName, + tapManager.Config.AgentImage, + TapperPodName, + fmt.Sprintf("%s.%s.svc.cluster.local", tapManager.state.apiServerService.Name, tapManager.state.apiServerService.Namespace), + nodeToTappedPodIPMap, + serviceAccountName, + tapManager.Config.TapperResources, + tapManager.Config.ImagePullPolicy, + 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 { + return err + } + } + + return nil +} \ No newline at end of file diff --git a/shared/kubernetes/provider.go b/shared/kubernetes/provider.go index 0b19991d5..e0c111eaf 100644 --- a/shared/kubernetes/provider.go +++ b/shared/kubernetes/provider.go @@ -510,7 +510,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, serializedMizuApiFilteringOptions string, 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 *interface{}, 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 { @@ -522,6 +522,11 @@ func (provider *Provider) ApplyMizuTapperDaemonSet(ctx context.Context, namespac return err } + mizuApiFilteringOptionsJsonStr, err := json.Marshal(mizuApiFilteringOptions) + if err != nil { + return err + } + mizuCmd := []string{ "./mizuagent", "-i", "any", @@ -546,7 +551,7 @@ func (provider *Provider) ApplyMizuTapperDaemonSet(ctx context.Context, namespac applyconfcore.EnvVar().WithName(shared.HostModeEnvVar).WithValue("1"), applyconfcore.EnvVar().WithName(shared.TappedAddressesPerNodeDictEnvVar).WithValue(string(nodeToTappedPodIPMapJsonStr)), applyconfcore.EnvVar().WithName(shared.GoGCEnvVar).WithValue("12800"), - applyconfcore.EnvVar().WithName(shared.MizuFilteringOptionsEnvVar).WithValue(serializedMizuApiFilteringOptions), + applyconfcore.EnvVar().WithName(shared.MizuFilteringOptionsEnvVar).WithValue(string(mizuApiFilteringOptionsJsonStr)), ) agentContainer.WithEnv( applyconfcore.EnvVar().WithName(shared.NodeNameEnvVar).WithValueFrom( diff --git a/shared/kubernetes/utils.go b/shared/kubernetes/utils.go new file mode 100644 index 000000000..a50adc3fd --- /dev/null +++ b/shared/kubernetes/utils.go @@ -0,0 +1,16 @@ +package kubernetes + +import core "k8s.io/api/core/v1" + +func GetNodeHostToTappedPodIpsMap(tappedPods []core.Pod) map[string][]string { + nodeToTappedPodIPMap := make(map[string][]string, 0) + for _, pod := range tappedPods { + existingList := nodeToTappedPodIPMap[pod.Spec.NodeName] + if existingList == nil { + nodeToTappedPodIPMap[pod.Spec.NodeName] = []string{pod.Status.PodIP} + } else { + nodeToTappedPodIPMap[pod.Spec.NodeName] = append(nodeToTappedPodIPMap[pod.Spec.NodeName], pod.Status.PodIP) + } + } + return nodeToTappedPodIPMap +} \ No newline at end of file