mirror of
https://github.com/kubeshark/kubeshark.git
synced 2026-09-01 00:57:17 +00:00
Update tapRunner.go, config.go, and 7 more files...
This commit is contained in:
+11
-13
@@ -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)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -84,6 +84,7 @@ func WriteConfig(config *ConfigStruct) error {
|
||||
}
|
||||
|
||||
type updateConfigStruct func(*ConfigStruct)
|
||||
|
||||
func UpdateConfig(updateConfigStruct updateConfigStruct) error {
|
||||
configFile, err := GetConfigWithDefaults()
|
||||
if err != nil {
|
||||
|
||||
@@ -32,4 +32,4 @@ func (config *LogsConfig) FilePath() string {
|
||||
}
|
||||
|
||||
return config.FileStr
|
||||
}
|
||||
}
|
||||
|
||||
@@ -14,4 +14,3 @@ const (
|
||||
ConfigMapName = MizuResourcesPrefix + "config"
|
||||
MinKubernetesServerVersion = "1.16.0"
|
||||
)
|
||||
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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)
|
||||
|
||||
|
||||
@@ -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
|
||||
|
||||
Reference in New Issue
Block a user