From 5e07924aca80e2fe2b7c274b7c9546043fd92c7a Mon Sep 17 00:00:00 2001 From: gadotroee <55343099+gadotroee@users.noreply.github.com> Date: Tue, 15 Feb 2022 17:13:08 +0200 Subject: [PATCH] Revert "Develop -> main (Release 27.0) (#812)" (#813) This reverts commit 77078e78d1d0e500bd16b34ab6aa1893ca2445a8. Co-authored-by: Igor Gov --- agent/pkg/api/main.go | 22 ++--- agent/pkg/api/socket_server_handlers.go | 6 +- agent/pkg/middlewares/cors.go | 2 +- agent/pkg/resolver/resolver.go | 34 ++++---- cli/README.md.TEMPLATE | 2 +- cli/apiserver/provider.go | 57 +++---------- cli/cmd/install.go | 18 +++-- cli/cmd/installRunner.go | 81 +++++++++++++++++++ cli/cmd/tapRunner.go | 1 - cli/cmd/viewRunner.go | 13 ++- docs/CHANGELOG.md | 73 ++++++++++++++++- shared/kubernetes/provider.go | 9 +-- shared/kubernetes/proxy.go | 11 +-- shared/models.go | 1 - tap/api/api.go | 18 ++--- tap/cleaner.go | 12 +-- tap/extensions/amqp/main.go | 9 +-- tap/extensions/http/Makefile | 2 +- tap/extensions/http/handlers.go | 16 ++-- tap/extensions/http/main.go | 20 ++--- tap/extensions/http/main_test.go | 17 +++- tap/extensions/http/matcher.go | 13 ++- tap/extensions/kafka/main.go | 15 ++-- tap/extensions/kafka/matcher.go | 14 ++-- tap/extensions/kafka/request.go | 5 +- tap/extensions/kafka/response.go | 5 +- tap/extensions/redis/handlers.go | 10 ++- tap/extensions/redis/main.go | 15 ++-- tap/extensions/redis/matcher.go | 13 ++- tap/passive_tapper.go | 1 - tap/tcp_reader.go | 3 +- tap/tcp_stream_factory.go | 14 ++-- .../Pages/TrafficPage/TrafficPage.tsx | 12 +-- 33 files changed, 308 insertions(+), 236 deletions(-) create mode 100644 cli/cmd/installRunner.go diff --git a/agent/pkg/api/main.go b/agent/pkg/api/main.go index 685ef2516..3886af658 100644 --- a/agent/pkg/api/main.go +++ b/agent/pkg/api/main.go @@ -118,8 +118,8 @@ func startReadingChannel(outputItems <-chan *tapApi.OutputChannelItem, extension for item := range outputItems { extension := extensionsMap[item.Protocol.Name] - resolvedSource, resolvedDestionation, namespace := resolveIP(item.ConnectionInfo) - mizuEntry := extension.Dissector.Analyze(item, resolvedSource, resolvedDestionation, namespace) + resolvedSource, resolvedDestionation := resolveIP(item.ConnectionInfo) + mizuEntry := extension.Dissector.Analyze(item, resolvedSource, resolvedDestionation) if extension.Protocol.Name == "http" { if !disableOASValidation { var httpPair tapApi.HTTPRequestResponsePair @@ -158,32 +158,26 @@ func startReadingChannel(outputItems <-chan *tapApi.OutputChannelItem, extension } } -func resolveIP(connectionInfo *tapApi.ConnectionInfo) (resolvedSource string, resolvedDestination string, namespace string) { +func resolveIP(connectionInfo *tapApi.ConnectionInfo) (resolvedSource string, resolvedDestination string) { if k8sResolver != nil { unresolvedSource := connectionInfo.ClientIP - resolvedSourceObject := k8sResolver.Resolve(unresolvedSource) - if resolvedSourceObject == nil { + resolvedSource = k8sResolver.Resolve(unresolvedSource) + if resolvedSource == "" { logger.Log.Debugf("Cannot find resolved name to source: %s", unresolvedSource) if os.Getenv("SKIP_NOT_RESOLVED_SOURCE") == "1" { return } - } else { - resolvedSource = resolvedSourceObject.FullAddress } - unresolvedDestination := fmt.Sprintf("%s:%s", connectionInfo.ServerIP, connectionInfo.ServerPort) - resolvedDestinationObject := k8sResolver.Resolve(unresolvedDestination) - if resolvedDestinationObject == nil { + resolvedDestination = k8sResolver.Resolve(unresolvedDestination) + if resolvedDestination == "" { logger.Log.Debugf("Cannot find resolved name to dest: %s", unresolvedDestination) if os.Getenv("SKIP_NOT_RESOLVED_DEST") == "1" { return } - } else { - resolvedDestination = resolvedDestinationObject.FullAddress - namespace = resolvedDestinationObject.Namespace } } - return resolvedSource, resolvedDestination, namespace + return resolvedSource, resolvedDestination } func CheckIsServiceIP(address string) bool { diff --git a/agent/pkg/api/socket_server_handlers.go b/agent/pkg/api/socket_server_handlers.go index bedde49a7..a22508362 100644 --- a/agent/pkg/api/socket_server_handlers.go +++ b/agent/pkg/api/socket_server_handlers.go @@ -104,9 +104,9 @@ func (h *RoutesEventHandlers) WebSocketMessage(_ int, message []byte) { } func handleTLSLink(outboundLinkMessage models.WebsocketOutboundLinkMessage) { - resolvedNameObject := k8sResolver.Resolve(outboundLinkMessage.Data.DstIP) - if resolvedNameObject != nil { - outboundLinkMessage.Data.DstIP = resolvedNameObject.FullAddress + resolvedName := k8sResolver.Resolve(outboundLinkMessage.Data.DstIP) + if resolvedName != "" { + outboundLinkMessage.Data.DstIP = resolvedName } else if outboundLinkMessage.Data.SuggestedResolvedName != "" { outboundLinkMessage.Data.DstIP = outboundLinkMessage.Data.SuggestedResolvedName } diff --git a/agent/pkg/middlewares/cors.go b/agent/pkg/middlewares/cors.go index e5d711ad9..04afb297a 100644 --- a/agent/pkg/middlewares/cors.go +++ b/agent/pkg/middlewares/cors.go @@ -7,7 +7,7 @@ func CORSMiddleware() gin.HandlerFunc { c.Writer.Header().Set("Access-Control-Allow-Origin", "*") c.Writer.Header().Set("Access-Control-Allow-Credentials", "true") c.Writer.Header().Set("Access-Control-Allow-Headers", "Content-Type, Content-Length, Accept-Encoding, X-CSRF-Token, Authorization, accept, origin, Cache-Control, X-Requested-With, x-session-token") - c.Writer.Header().Set("Access-Control-Allow-Methods", "POST, OPTIONS, GET, PUT, DELETE") + c.Writer.Header().Set("Access-Control-Allow-Methods", "POST, OPTIONS, GET, PUT") if c.Request.Method == "OPTIONS" { c.AbortWithStatus(204) diff --git a/agent/pkg/resolver/resolver.go b/agent/pkg/resolver/resolver.go index a1cdc52de..60533704d 100644 --- a/agent/pkg/resolver/resolver.go +++ b/agent/pkg/resolver/resolver.go @@ -30,11 +30,6 @@ type Resolver struct { namespace string } -type ResolvedObjectInfo struct { - FullAddress string - Namespace string -} - func (resolver *Resolver) Start(ctx context.Context) { if !resolver.isStarted { resolver.isStarted = true @@ -45,12 +40,12 @@ func (resolver *Resolver) Start(ctx context.Context) { } } -func (resolver *Resolver) Resolve(name string) *ResolvedObjectInfo { +func (resolver *Resolver) Resolve(name string) string { resolvedName, isFound := resolver.nameMap.Get(name) if !isFound { - return nil + return "" } - return resolvedName.(*ResolvedObjectInfo) + return resolvedName.(string) } func (resolver *Resolver) GetMap() cmap.ConcurrentMap { @@ -76,7 +71,7 @@ func (resolver *Resolver) watchPods(ctx context.Context) error { } if event.Type == watch.Deleted { pod := event.Object.(*corev1.Pod) - resolver.saveResolvedName(pod.Status.PodIP, "", pod.Namespace, event.Type) + resolver.saveResolvedName(pod.Status.PodIP, "", event.Type) } case <-ctx.Done(): watcher.Stop() @@ -111,10 +106,10 @@ func (resolver *Resolver) watchEndpoints(ctx context.Context) error { } if subset.Addresses != nil { for _, address := range subset.Addresses { - resolver.saveResolvedName(address.IP, serviceHostname, endpoint.Namespace, event.Type) + resolver.saveResolvedName(address.IP, serviceHostname, event.Type) for _, port := range ports { ipWithPort := fmt.Sprintf("%s:%d", address.IP, port) - resolver.saveResolvedName(ipWithPort, serviceHostname, endpoint.Namespace, event.Type) + resolver.saveResolvedName(ipWithPort, serviceHostname, event.Type) } } } @@ -144,19 +139,19 @@ func (resolver *Resolver) watchServices(ctx context.Context) error { service := event.Object.(*corev1.Service) serviceHostname := fmt.Sprintf("%s.%s", service.Name, service.Namespace) if service.Spec.ClusterIP != "" && service.Spec.ClusterIP != kubClientNullString { - resolver.saveResolvedName(service.Spec.ClusterIP, serviceHostname, service.Namespace, event.Type) + resolver.saveResolvedName(service.Spec.ClusterIP, serviceHostname, event.Type) if service.Spec.Ports != nil { for _, port := range service.Spec.Ports { if port.Port > 0 { - resolver.saveResolvedName(fmt.Sprintf("%s:%d", service.Spec.ClusterIP, port.Port), serviceHostname, service.Namespace, event.Type) + resolver.saveResolvedName(fmt.Sprintf("%s:%d", service.Spec.ClusterIP, port.Port), serviceHostname, event.Type) } } } - resolver.saveServiceIP(service.Spec.ClusterIP, serviceHostname, service.Namespace, event.Type) + resolver.saveServiceIP(service.Spec.ClusterIP, serviceHostname, event.Type) } if service.Status.LoadBalancer.Ingress != nil { for _, ingress := range service.Status.LoadBalancer.Ingress { - resolver.saveResolvedName(ingress.IP, serviceHostname, service.Namespace, event.Type) + resolver.saveResolvedName(ingress.IP, serviceHostname, event.Type) } } case <-ctx.Done(): @@ -166,22 +161,21 @@ func (resolver *Resolver) watchServices(ctx context.Context) error { } } -func (resolver *Resolver) saveResolvedName(key string, resolved string, namespace string, eventType watch.EventType) { +func (resolver *Resolver) saveResolvedName(key string, resolved string, eventType watch.EventType) { if eventType == watch.Deleted { resolver.nameMap.Remove(key) logger.Log.Infof("setting %s=nil", key) } else { - - resolver.nameMap.Set(key, &ResolvedObjectInfo{FullAddress: resolved, Namespace: namespace}) + resolver.nameMap.Set(key, resolved) logger.Log.Infof("setting %s=%s", key, resolved) } } -func (resolver *Resolver) saveServiceIP(key string, resolved string, namespace string, eventType watch.EventType) { +func (resolver *Resolver) saveServiceIP(key string, resolved string, eventType watch.EventType) { if eventType == watch.Deleted { resolver.serviceMap.Remove(key) } else { - resolver.nameMap.Set(key, &ResolvedObjectInfo{FullAddress: resolved, Namespace: namespace}) + resolver.serviceMap.Set(key, resolved) } } diff --git a/cli/README.md.TEMPLATE b/cli/README.md.TEMPLATE index 0271bdc27..fee03a253 100644 --- a/cli/README.md.TEMPLATE +++ b/cli/README.md.TEMPLATE @@ -1,5 +1,5 @@ # Mizu release _VER_ -Mizu CHANGELOG is now part of [Mizu wiki](https://github.com/up9inc/mizu/wiki/CHANGELOG) +Full changelog for stable release see in [docs](https://github.com/up9inc/mizu/blob/main/docs/CHANGELOG.md) ## Download Mizu for your platform diff --git a/cli/apiserver/provider.go b/cli/apiserver/provider.go index b4fd4b394..07afe6803 100644 --- a/cli/apiserver/provider.go +++ b/cli/apiserver/provider.go @@ -4,11 +4,9 @@ import ( "bytes" "encoding/json" "fmt" - "io" "io/ioutil" "net/http" "net/url" - "strings" "time" "github.com/up9inc/mizu/shared/kubernetes" @@ -59,8 +57,10 @@ func (provider *Provider) TestConnection() error { func (provider *Provider) isReachable() (bool, error) { echoUrl := fmt.Sprintf("%s/echo", provider.url) - if _, err := provider.get(echoUrl); err != nil { + if response, err := provider.client.Get(echoUrl); err != nil { return false, err + } else if response.StatusCode != 200 { + return false, fmt.Errorf("invalid status code %v", response.StatusCode) } else { return true, nil } @@ -72,8 +72,10 @@ func (provider *Provider) ReportTapperStatus(tapperStatus shared.TapperStatus) e if jsonValue, err := json.Marshal(tapperStatus); err != nil { return fmt.Errorf("failed Marshal the tapper status %w", err) } else { - if _, err := provider.post(tapperStatusUrl, "application/json", bytes.NewBuffer(jsonValue)); err != nil { + if response, err := provider.client.Post(tapperStatusUrl, "application/json", bytes.NewBuffer(jsonValue)); err != nil { return fmt.Errorf("failed sending to API server the tapped pods %w", err) + } else if response.StatusCode != 200 { + return fmt.Errorf("failed sending to API server the tapper status, response status code %v", response.StatusCode) } else { logger.Log.Debugf("Reported to server API about tapper status: %v", tapperStatus) return nil @@ -89,8 +91,10 @@ func (provider *Provider) ReportTappedPods(pods []core.Pod) error { if jsonValue, err := json.Marshal(podInfos); err != nil { return fmt.Errorf("failed Marshal the tapped pods %w", err) } else { - if _, err := provider.post(tappedPodsUrl, "application/json", bytes.NewBuffer(jsonValue)); err != nil { + if response, err := provider.client.Post(tappedPodsUrl, "application/json", bytes.NewBuffer(jsonValue)); err != nil { return fmt.Errorf("failed sending to API server the tapped pods %w", err) + } else if response.StatusCode != 200 { + return fmt.Errorf("failed sending to API server the tapped pods, response status code %v", response.StatusCode) } else { logger.Log.Debugf("Reported to server API about %d taped pods successfully", len(podInfos)) return nil @@ -101,9 +105,11 @@ func (provider *Provider) ReportTappedPods(pods []core.Pod) error { func (provider *Provider) GetGeneralStats() (map[string]interface{}, error) { generalStatsUrl := fmt.Sprintf("%s/status/general", provider.url) - response, requestErr := provider.get(generalStatsUrl) + response, requestErr := provider.client.Get(generalStatsUrl) if requestErr != nil { return nil, fmt.Errorf("failed to get general stats for telemetry, err: %w", requestErr) + } else if response.StatusCode != 200 { + return nil, fmt.Errorf("failed to get general stats for telemetry, status code: %v", response.StatusCode) } defer response.Body.Close() @@ -126,7 +132,7 @@ func (provider *Provider) GetVersion() (string, error) { Method: http.MethodGet, URL: versionUrl, } - statusResp, err := provider.do(req) + statusResp, err := provider.client.Do(req) if err != nil { return "", err } @@ -139,40 +145,3 @@ func (provider *Provider) GetVersion() (string, error) { return versionResponse.Ver, nil } - -// When err is nil, resp always contains a non-nil resp.Body. -// Caller should close resp.Body when done reading from it. -func (provider *Provider) get(url string) (*http.Response, error) { - return provider.checkError(provider.client.Get(url)) -} - -// When err is nil, resp always contains a non-nil resp.Body. -// Caller should close resp.Body when done reading from it. -func (provider *Provider) post(url, contentType string, body io.Reader) (*http.Response, error) { - return provider.checkError(provider.client.Post(url, contentType, body)) -} - -// When err is nil, resp always contains a non-nil resp.Body. -// Caller should close resp.Body when done reading from it. -func (provider *Provider) do(req *http.Request) (*http.Response, error) { - return provider.checkError(provider.client.Do(req)) -} - -func (provider *Provider) checkError(response *http.Response, errInOperation error) (*http.Response, error) { - if (errInOperation != nil) { - return response, errInOperation - // Check only if status != 200 (and not status >= 300). Agent APIs return only 200 on success. - } else if response.StatusCode != http.StatusOK { - body, err := ioutil.ReadAll(response.Body) - response.Body.Close() - response.Body = io.NopCloser(bytes.NewBuffer(body)) // rewind - if err != nil { - return response, err - } - - errorMsg := strings.ReplaceAll((string(body)), "\n", ";") - return response, fmt.Errorf("got response with status code: %d, body: %s", response.StatusCode, errorMsg) - } - - return response, nil -} diff --git a/cli/cmd/install.go b/cli/cmd/install.go index 509538d61..e17b1258b 100644 --- a/cli/cmd/install.go +++ b/cli/cmd/install.go @@ -1,9 +1,10 @@ package cmd import ( + "fmt" "github.com/spf13/cobra" + "github.com/up9inc/mizu/cli/config" "github.com/up9inc/mizu/cli/telemetry" - "github.com/up9inc/mizu/shared/logger" ) var installCmd = &cobra.Command{ @@ -11,13 +12,13 @@ var installCmd = &cobra.Command{ Short: "Installs mizu components", RunE: func(cmd *cobra.Command, args []string) error { go telemetry.ReportRun("install", nil) - logger.Log.Infof("This command has been deprecated, please use helm as described below.\n\n") - - logger.Log.Infof("To install stable build of Mizu on your cluster using helm, run the following command:") - logger.Log.Infof(" helm install mizu https://static.up9.com/mizu/helm --namespace=mizu-ent --create-namespace\n\n") - - logger.Log.Infof("To install development build of Mizu on your cluster using helm, run the following command:") - logger.Log.Infof(" helm install mizu https://static.up9.com/mizu/helm-develop --namespace=mizu-ent --create-namespace") + runMizuInstall() + return nil + }, + PreRunE: func(cmd *cobra.Command, args []string) error { + if config.Config.IsNsRestrictedMode() { + return fmt.Errorf("install is not supported in restricted namespace mode") + } return nil }, @@ -26,3 +27,4 @@ var installCmd = &cobra.Command{ func init() { rootCmd.AddCommand(installCmd) } + diff --git a/cli/cmd/installRunner.go b/cli/cmd/installRunner.go new file mode 100644 index 000000000..840e5710e --- /dev/null +++ b/cli/cmd/installRunner.go @@ -0,0 +1,81 @@ +package cmd + +import ( + "context" + "errors" + "fmt" + + "github.com/creasty/defaults" + "github.com/up9inc/mizu/cli/config" + "github.com/up9inc/mizu/cli/errormessage" + "github.com/up9inc/mizu/cli/resources" + "github.com/up9inc/mizu/cli/uiUtils" + "github.com/up9inc/mizu/shared" + "github.com/up9inc/mizu/shared/logger" + k8serrors "k8s.io/apimachinery/pkg/api/errors" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" +) + +func runMizuInstall() { + kubernetesProvider, err := getKubernetesProviderForCli() + if err != nil { + return + } + + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() // cancel will be called when this function exits + + var serializedValidationRules string + var serializedContract string + + var defaultMaxEntriesDBSizeBytes int64 = 200 * 1000 * 1000 + + defaultResources := shared.Resources{} + if err := defaults.Set(&defaultResources); err != nil { + logger.Log.Debug(err) + } + + mizuAgentConfig := getInstallMizuAgentConfig(defaultMaxEntriesDBSizeBytes, defaultResources) + serializedMizuConfig, err := getSerializedMizuAgentConfig(mizuAgentConfig) + if err != nil { + logger.Log.Errorf(uiUtils.Error, fmt.Sprintf("Error serializing mizu config: %v", errormessage.FormatError(err))) + return + } + + if err = resources.CreateInstallMizuResources(ctx, kubernetesProvider, serializedValidationRules, + serializedContract, serializedMizuConfig, config.Config.IsNsRestrictedMode(), + config.Config.MizuResourcesNamespace, config.Config.AgentImage, + config.Config.KratosImage, config.Config.KetoImage, + nil, defaultMaxEntriesDBSizeBytes, defaultResources, config.Config.ImagePullPolicy(), + config.Config.LogLevel(), false); err != nil { + var statusError *k8serrors.StatusError + if errors.As(err, &statusError) && (statusError.ErrStatus.Reason == metav1.StatusReasonAlreadyExists) { + logger.Log.Info("Mizu is already running in this namespace, run `mizu clean` to remove the currently running Mizu instance") + } else { + defer resources.CleanUpMizuResources(ctx, cancel, kubernetesProvider, config.Config.IsNsRestrictedMode(), config.Config.MizuResourcesNamespace) + logger.Log.Errorf(uiUtils.Error, fmt.Sprintf("Error creating resources: %v", errormessage.FormatError(err))) + } + + return + } + + logger.Log.Infof(uiUtils.Magenta, "Installation completed, run `mizu view` to connect to the mizu daemon instance") +} + +func getInstallMizuAgentConfig(maxDBSizeBytes int64, tapperResources shared.Resources) *shared.MizuAgentConfig { + mizuAgentConfig := shared.MizuAgentConfig{ + MaxDBSizeBytes: maxDBSizeBytes, + AgentImage: config.Config.AgentImage, + PullPolicy: config.Config.ImagePullPolicyStr, + LogLevel: config.Config.LogLevel(), + TapperResources: tapperResources, + MizuResourcesNamespace: config.Config.MizuResourcesNamespace, + AgentDatabasePath: shared.DataDirPath, + StandaloneMode: true, + ServiceMap: config.Config.ServiceMap, + OAS: config.Config.OAS, + Elastic: config.Config.Elastic, + } + + return &mizuAgentConfig +} diff --git a/cli/cmd/tapRunner.go b/cli/cmd/tapRunner.go index 4e8383446..5e37135bc 100644 --- a/cli/cmd/tapRunner.go +++ b/cli/cmd/tapRunner.go @@ -162,7 +162,6 @@ func getTapMizuAgentConfig() *shared.MizuAgentConfig { AgentDatabasePath: shared.DataDirPath, ServiceMap: config.Config.ServiceMap, OAS: config.Config.OAS, - Telemetry: config.Config.Telemetry, Elastic: config.Config.Elastic, } diff --git a/cli/cmd/viewRunner.go b/cli/cmd/viewRunner.go index ebcddc67e..9a4101edd 100644 --- a/cli/cmd/viewRunner.go +++ b/cli/cmd/viewRunner.go @@ -3,13 +3,13 @@ package cmd import ( "context" "fmt" - "net/http" - "github.com/up9inc/mizu/cli/utils" + "net/http" "github.com/up9inc/mizu/cli/apiserver" "github.com/up9inc/mizu/cli/config" "github.com/up9inc/mizu/cli/mizu/fsUtils" + "github.com/up9inc/mizu/cli/mizu/version" "github.com/up9inc/mizu/cli/uiUtils" "github.com/up9inc/mizu/shared/kubernetes" "github.com/up9inc/mizu/shared/logger" @@ -62,5 +62,14 @@ func runMizuView() { uiUtils.OpenBrowser(url) } + if isCompatible, err := version.CheckVersionCompatibility(apiServerProvider); err != nil { + logger.Log.Errorf("Failed to check versions compatibility %v", err) + cancel() + return + } else if !isCompatible { + cancel() + return + } + utils.WaitForFinish(ctx, cancel) } diff --git a/docs/CHANGELOG.md b/docs/CHANGELOG.md index c27c061b0..a866207cf 100644 --- a/docs/CHANGELOG.md +++ b/docs/CHANGELOG.md @@ -1 +1,72 @@ -Mizu CHANGELOG is now part of [Mizu wiki](https://github.com/up9inc/mizu/wiki/CHANGELOG) +# CHANGELOG +This document summarizes main and fixes changes published in stable (aka `main`) branch of this project. +Ongoing work and development releases are under `develop` branch. + +## 0.24.0 + +### main features +* ARM64 support -- Mizu is now available for ARM 64bit architecture + * Now you can run Mizu with `minikube` on your Apple M1 laptop or any other ARM-based hosts +* New command helps user verify Mizu deployment + * Run `mizu check` to verify Mizu was deployed successfully + * `mizu check` verifies version compatibility, resources and permissions required by Mizu +* EXPERIMENTAL: Service Map - graph of all service interactions + * Arrow direction show client to server connection + * Graph edge width reflects volume of traffic captured between the services + * to enable this experimental feature use `--set service-map=true` flag + +### improvements +* Mizu container images are now served from [Docker Hub](https://hub.docker.com/r/up9inc/mizu), as multi-architecture images (arm64, amd64) +* in Mizu GUI the filter query can now be applied by pressing CONTROL/COMMAND + ENTER +* try port-forwarding if http-proxy connection to Mizu API server is not available + +### notable bug fixes +* Fixed HTTP/1.0 presentation which was shown as HTTP/1.1 +* Fixed handling of long-living TCP connections, improves capturing gRPC and HTTP/2 traffic, and helps in service-mesh setups (istio, linkerd) + + +## 0.23.0 +### notable bug fixes +* fixed errors in Redis protocol parser (better handling of Array and Bulk String message types) + + + +## 0.22.0 + +### main features +* Service Mesh support -- mizu is now capable to tap mTLS traffic between pods connected by Istio service mesh + * Use `--service-mesh` option to enable this feature +* New installation option - have the same Mizu functionality as long living pods in your cluster, with password protection + * To install use `mizu install` command + * To access use `mizu view` or `kubectl -n mizu port-forward svc/mizu-api-server` + * To uninstall run `mizu clean` +* At first login + * Set admin password as prompted, use it to login to mizu later on. + * After login, user should select cluster namespaces to tap: by default all namespaces in the cluster are selected, user can select/unselect according to their needs. These settings are retained and can be modified at any time via Settings menu (cog icon on the top-right) + + +### improvements +* improved Mizu permissions/roles logic to support clusters with strict PodSecurityPolicy (PSP) -- see [PERMISSIONS](PERMISSIONS.md) doc for more details + +### notable bug fixes +* mizu now works properly when API service is exposed via HTTPS url +* mizu now properly displays KAFKA message body + + + + +## 0.21.0 + +### main features +* New traffic search & stream exprience +* Rich query language with full-text search capabilities on headers & body +* Distinct live-streaming vs paging/browsing modes, all with filter applied + +### improvements +* GUI - source and destination IP addresses & service names for each traffic item +* GUI - Mizu health - display warning sign in top bar when not all requested pods are successfully tapped +* GUI - pod tapping status reflected in the list (ok or problem) +* Mizu telemetry - report platform type + +### fixes +* Request duration and body size properly shown in GUI (instead of -1) diff --git a/shared/kubernetes/provider.go b/shared/kubernetes/provider.go index 3ce71a56b..fbb74df05 100644 --- a/shared/kubernetes/provider.go +++ b/shared/kubernetes/provider.go @@ -76,8 +76,6 @@ func NewProvider(kubeConfigPath string) (*Provider, error) { "you can set alternative kube config file path by adding the kube-config-path field to the mizu config file, err: %w", kubeConfigPath, err) } - logger.Log.Debugf("K8s client config, host: %s, api path: %s, user agent: %s", restClientConfig.Host, restClientConfig.APIPath, restClientConfig.UserAgent) - return &Provider{ clientSet: clientSet, kubernetesConfig: kubernetesConfig, @@ -954,11 +952,6 @@ func (provider *Provider) ApplyMizuTapperDaemonSet(ctx context.Context, namespac labelSelector := applyconfmeta.LabelSelector() labelSelector.WithMatchLabels(map[string]string{"app": tapperPodName}) - applyOptions := metav1.ApplyOptions{ - Force: true, - FieldManager: fieldManagerName, - } - daemonSet := applyconfapp.DaemonSet(daemonSetName, namespace) daemonSet. WithLabels(map[string]string{ @@ -967,7 +960,7 @@ func (provider *Provider) ApplyMizuTapperDaemonSet(ctx context.Context, namespac }). WithSpec(applyconfapp.DaemonSetSpec().WithSelector(labelSelector).WithTemplate(podTemplate)) - _, err = provider.clientSet.AppsV1().DaemonSets(namespace).Apply(ctx, daemonSet, applyOptions) + _, err = provider.clientSet.AppsV1().DaemonSets(namespace).Apply(ctx, daemonSet, metav1.ApplyOptions{FieldManager: fieldManagerName}) return err } diff --git a/shared/kubernetes/proxy.go b/shared/kubernetes/proxy.go index 521a90b57..b21d7d097 100644 --- a/shared/kubernetes/proxy.go +++ b/shared/kubernetes/proxy.go @@ -128,14 +128,9 @@ func getHttpDialer(kubernetesProvider *Provider, namespace string, podName strin return nil, err } - clientConfigHostUrl, err := url.Parse(kubernetesProvider.clientConfig.Host) - if err != nil { - return nil, fmt.Errorf("Failed parsing client config host URL %s, error %w", kubernetesProvider.clientConfig.Host, err) - } - path := fmt.Sprintf("%s/api/v1/namespaces/%s/pods/%s/portforward", clientConfigHostUrl.Path, namespace, podName) - - serverURL := url.URL{Scheme: "https", Path: path, Host: clientConfigHostUrl.Host} - logger.Log.Debugf("Http dialer url %v", serverURL) + path := fmt.Sprintf("/api/v1/namespaces/%s/pods/%s/portforward", namespace, podName) + hostIP := strings.TrimLeft(kubernetesProvider.clientConfig.Host, "htps:/") // no need specify "t" twice + serverURL := url.URL{Scheme: "https", Path: path, Host: hostIP} return spdy.NewDialer(upgrader, &http.Client{Transport: roundTripper}, http.MethodPost, &serverURL), nil } diff --git a/shared/models.go b/shared/models.go index 7aa0d3d47..9081e9a9e 100644 --- a/shared/models.go +++ b/shared/models.go @@ -43,7 +43,6 @@ type MizuAgentConfig struct { StandaloneMode bool `json:"standaloneMode"` ServiceMap bool `json:"serviceMap"` OAS bool `json:"oas"` - Telemetry bool `json:"telemetry"` Elastic ElasticConfig `json:"elastic"` } diff --git a/tap/api/api.go b/tap/api/api.go index d1a812e80..52c3cac11 100644 --- a/tap/api/api.go +++ b/tap/api/api.go @@ -39,9 +39,10 @@ type TCP struct { } type Extension struct { - Protocol *Protocol - Path string - Dissector Dissector + Protocol *Protocol + Path string + Dissector Dissector + MatcherMap *sync.Map } type ConnectionInfo struct { @@ -61,6 +62,7 @@ type TcpID struct { } type CounterPair struct { + StreamId int64 Request uint Response uint sync.Mutex @@ -98,15 +100,10 @@ type SuperIdentifier struct { type Dissector interface { Register(*Extension) Ping() - Dissect(b *bufio.Reader, isClient bool, tcpID *TcpID, counterPair *CounterPair, superTimer *SuperTimer, superIdentifier *SuperIdentifier, emitter Emitter, options *TrafficFilteringOptions, reqResMatcher RequestResponseMatcher) error - Analyze(item *OutputChannelItem, resolvedSource string, resolvedDestination string, namespace string) *Entry + Dissect(b *bufio.Reader, isClient bool, tcpID *TcpID, counterPair *CounterPair, superTimer *SuperTimer, superIdentifier *SuperIdentifier, emitter Emitter, options *TrafficFilteringOptions) error + Analyze(item *OutputChannelItem, resolvedSource string, resolvedDestination string) *Entry Represent(request map[string]interface{}, response map[string]interface{}) (object []byte, bodySize int64, err error) Macros() map[string]string - NewResponseRequestMatcher() RequestResponseMatcher -} - -type RequestResponseMatcher interface { - GetMap() *sync.Map } type Emitting struct { @@ -128,7 +125,6 @@ type Entry struct { Protocol Protocol `json:"proto"` Source *TCP `json:"src"` Destination *TCP `json:"dst"` - Namespace string `json:"namespace,omitempty"` Outgoing bool `json:"outgoing"` Timestamp int64 `json:"timestamp"` StartTime time.Time `json:"startTime"` diff --git a/tap/cleaner.go b/tap/cleaner.go index cdf2acf20..61be717f3 100644 --- a/tap/cleaner.go +++ b/tap/cleaner.go @@ -22,7 +22,6 @@ type Cleaner struct { connectionTimeout time.Duration stats CleanerStats statsMutex sync.Mutex - streamsMap *tcpStreamMap } func (cl *Cleaner) clean() { @@ -33,15 +32,10 @@ func (cl *Cleaner) clean() { flushed, closed := cl.assembler.FlushCloseOlderThan(startCleanTime.Add(-cl.connectionTimeout)) cl.assemblerMutex.Unlock() - cl.streamsMap.streams.Range(func(k, v interface{}) bool { - reqResMatcher := v.(*tcpStreamWrapper).reqResMatcher - if reqResMatcher == nil { - return true - } - deleted := deleteOlderThan(reqResMatcher.GetMap(), startCleanTime.Add(-cl.connectionTimeout)) + for _, extension := range extensions { + deleted := deleteOlderThan(extension.MatcherMap, startCleanTime.Add(-cl.connectionTimeout)) cl.stats.deleted += deleted - return true - }) + } cl.statsMutex.Lock() logger.Log.Debugf("Assembler Stats after cleaning %s", cl.assembler.Dump()) diff --git a/tap/extensions/amqp/main.go b/tap/extensions/amqp/main.go index f672dba7b..9ca023a4f 100644 --- a/tap/extensions/amqp/main.go +++ b/tap/extensions/amqp/main.go @@ -42,7 +42,7 @@ func (d dissecting) Ping() { const amqpRequest string = "amqp_request" -func (d dissecting) Dissect(b *bufio.Reader, isClient bool, tcpID *api.TcpID, counterPair *api.CounterPair, superTimer *api.SuperTimer, superIdentifier *api.SuperIdentifier, emitter api.Emitter, options *api.TrafficFilteringOptions, _reqResMatcher api.RequestResponseMatcher) error { +func (d dissecting) Dissect(b *bufio.Reader, isClient bool, tcpID *api.TcpID, counterPair *api.CounterPair, superTimer *api.SuperTimer, superIdentifier *api.SuperIdentifier, emitter api.Emitter, options *api.TrafficFilteringOptions) error { r := AmqpReader{b} var remaining int @@ -212,7 +212,7 @@ func (d dissecting) Dissect(b *bufio.Reader, isClient bool, tcpID *api.TcpID, co } } -func (d dissecting) Analyze(item *api.OutputChannelItem, resolvedSource string, resolvedDestination string, namespace string) *api.Entry { +func (d dissecting) Analyze(item *api.OutputChannelItem, resolvedSource string, resolvedDestination string) *api.Entry { request := item.Pair.Request.Payload.(map[string]interface{}) reqDetails := request["details"].(map[string]interface{}) @@ -254,7 +254,6 @@ func (d dissecting) Analyze(item *api.OutputChannelItem, resolvedSource string, IP: item.ConnectionInfo.ServerIP, Port: item.ConnectionInfo.ServerPort, }, - Namespace: namespace, Outgoing: item.ConnectionInfo.IsOutgoing, Request: reqDetails, Method: request["method"].(string), @@ -301,10 +300,6 @@ func (d dissecting) Macros() map[string]string { } } -func (d dissecting) NewResponseRequestMatcher() api.RequestResponseMatcher { - return nil -} - var Dissector dissecting func NewDissector() api.Dissector { diff --git a/tap/extensions/http/Makefile b/tap/extensions/http/Makefile index 253910d58..529cc27ef 100644 --- a/tap/extensions/http/Makefile +++ b/tap/extensions/http/Makefile @@ -13,4 +13,4 @@ test-pull-bin: test-pull-expect: @mkdir -p expect - @[ "${skipexpect}" ] && echo "Skipping downloading expected JSONs" || gsutil -o 'GSUtil:parallel_process_count=5' -o 'GSUtil:parallel_thread_count=5' -m cp -r gs://static.up9.io/mizu/test-pcap/expect2/http/\* expect + @[ "${skipexpect}" ] && echo "Skipping downloading expected JSONs" || gsutil -o 'GSUtil:parallel_process_count=5' -o 'GSUtil:parallel_thread_count=5' -m cp -r gs://static.up9.io/mizu/test-pcap/expect/http/\* expect diff --git a/tap/extensions/http/handlers.go b/tap/extensions/http/handlers.go index 7ffa5fee5..8d084be77 100644 --- a/tap/extensions/http/handlers.go +++ b/tap/extensions/http/handlers.go @@ -47,7 +47,7 @@ func replaceForwardedFor(item *api.OutputChannelItem) { item.ConnectionInfo.ClientPort = "" } -func handleHTTP2Stream(http2Assembler *Http2Assembler, tcpID *api.TcpID, superTimer *api.SuperTimer, emitter api.Emitter, options *api.TrafficFilteringOptions, reqResMatcher *requestResponseMatcher) error { +func handleHTTP2Stream(http2Assembler *Http2Assembler, tcpID *api.TcpID, superTimer *api.SuperTimer, emitter api.Emitter, options *api.TrafficFilteringOptions) error { streamID, messageHTTP1, isGrpc, err := http2Assembler.readMessage() if err != nil { return err @@ -58,7 +58,7 @@ func handleHTTP2Stream(http2Assembler *Http2Assembler, tcpID *api.TcpID, superTi switch messageHTTP1 := messageHTTP1.(type) { case http.Request: ident := fmt.Sprintf( - "%s_%s_%s_%s_%d_%s", + "%s->%s %s->%s %d %s", tcpID.SrcIP, tcpID.DstIP, tcpID.SrcPort, @@ -78,7 +78,7 @@ func handleHTTP2Stream(http2Assembler *Http2Assembler, tcpID *api.TcpID, superTi } case http.Response: ident := fmt.Sprintf( - "%s_%s_%s_%s_%d_%s", + "%s->%s %s->%s %d %s", tcpID.DstIP, tcpID.SrcIP, tcpID.DstPort, @@ -110,7 +110,7 @@ func handleHTTP2Stream(http2Assembler *Http2Assembler, tcpID *api.TcpID, superTi return nil } -func handleHTTP1ClientStream(b *bufio.Reader, tcpID *api.TcpID, counterPair *api.CounterPair, superTimer *api.SuperTimer, emitter api.Emitter, options *api.TrafficFilteringOptions, reqResMatcher *requestResponseMatcher) (switchingProtocolsHTTP2 bool, req *http.Request, err error) { +func handleHTTP1ClientStream(b *bufio.Reader, tcpID *api.TcpID, counterPair *api.CounterPair, superTimer *api.SuperTimer, emitter api.Emitter, options *api.TrafficFilteringOptions) (switchingProtocolsHTTP2 bool, req *http.Request, err error) { req, err = http.ReadRequest(b) if err != nil { return @@ -130,7 +130,8 @@ func handleHTTP1ClientStream(b *bufio.Reader, tcpID *api.TcpID, counterPair *api req.Body = io.NopCloser(bytes.NewBuffer(body)) // rewind ident := fmt.Sprintf( - "%s_%s_%s_%s_%d_%s", + "%d_%s:%s_%s:%s_%d_%s", + counterPair.StreamId, tcpID.SrcIP, tcpID.DstIP, tcpID.SrcPort, @@ -152,7 +153,7 @@ func handleHTTP1ClientStream(b *bufio.Reader, tcpID *api.TcpID, counterPair *api return } -func handleHTTP1ServerStream(b *bufio.Reader, tcpID *api.TcpID, counterPair *api.CounterPair, superTimer *api.SuperTimer, emitter api.Emitter, options *api.TrafficFilteringOptions, reqResMatcher *requestResponseMatcher) (switchingProtocolsHTTP2 bool, err error) { +func handleHTTP1ServerStream(b *bufio.Reader, tcpID *api.TcpID, counterPair *api.CounterPair, superTimer *api.SuperTimer, emitter api.Emitter, options *api.TrafficFilteringOptions) (switchingProtocolsHTTP2 bool, err error) { var res *http.Response res, err = http.ReadResponse(b, nil) if err != nil { @@ -173,7 +174,8 @@ func handleHTTP1ServerStream(b *bufio.Reader, tcpID *api.TcpID, counterPair *api res.Body = io.NopCloser(bytes.NewBuffer(body)) // rewind ident := fmt.Sprintf( - "%s_%s_%s_%s_%d_%s", + "%d_%s:%s_%s:%s_%d_%s", + counterPair.StreamId, tcpID.DstIP, tcpID.SrcIP, tcpID.DstPort, diff --git a/tap/extensions/http/main.go b/tap/extensions/http/main.go index 2b3d781f1..28c8a1744 100644 --- a/tap/extensions/http/main.go +++ b/tap/extensions/http/main.go @@ -84,15 +84,14 @@ type dissecting string func (d dissecting) Register(extension *api.Extension) { extension.Protocol = &http11protocol + extension.MatcherMap = reqResMatcher.openMessagesMap } func (d dissecting) Ping() { log.Printf("pong %s", http11protocol.Name) } -func (d dissecting) Dissect(b *bufio.Reader, isClient bool, tcpID *api.TcpID, counterPair *api.CounterPair, superTimer *api.SuperTimer, superIdentifier *api.SuperIdentifier, emitter api.Emitter, options *api.TrafficFilteringOptions, _reqResMatcher api.RequestResponseMatcher) error { - reqResMatcher := _reqResMatcher.(*requestResponseMatcher) - +func (d dissecting) Dissect(b *bufio.Reader, isClient bool, tcpID *api.TcpID, counterPair *api.CounterPair, superTimer *api.SuperTimer, superIdentifier *api.SuperIdentifier, emitter api.Emitter, options *api.TrafficFilteringOptions) error { var err error isHTTP2, _ := checkIsHTTP2Connection(b, isClient) @@ -125,7 +124,7 @@ func (d dissecting) Dissect(b *bufio.Reader, isClient bool, tcpID *api.TcpID, co } if isHTTP2 { - err = handleHTTP2Stream(http2Assembler, tcpID, superTimer, emitter, options, reqResMatcher) + err = handleHTTP2Stream(http2Assembler, tcpID, superTimer, emitter, options) if err == io.EOF || err == io.ErrUnexpectedEOF { break } else if err != nil { @@ -134,7 +133,7 @@ func (d dissecting) Dissect(b *bufio.Reader, isClient bool, tcpID *api.TcpID, co superIdentifier.Protocol = &http11protocol } else if isClient { var req *http.Request - switchingProtocolsHTTP2, req, err = handleHTTP1ClientStream(b, tcpID, counterPair, superTimer, emitter, options, reqResMatcher) + switchingProtocolsHTTP2, req, err = handleHTTP1ClientStream(b, tcpID, counterPair, superTimer, emitter, options) if err == io.EOF || err == io.ErrUnexpectedEOF { break } else if err != nil { @@ -145,7 +144,7 @@ func (d dissecting) Dissect(b *bufio.Reader, isClient bool, tcpID *api.TcpID, co // In case of an HTTP2 upgrade, duplicate the HTTP1 request into HTTP2 with stream ID 1 if switchingProtocolsHTTP2 { ident := fmt.Sprintf( - "%s_%s_%s_%s_1_%s", + "%s->%s %s->%s 1 %s", tcpID.SrcIP, tcpID.DstIP, tcpID.SrcPort, @@ -165,7 +164,7 @@ func (d dissecting) Dissect(b *bufio.Reader, isClient bool, tcpID *api.TcpID, co } } } else { - switchingProtocolsHTTP2, err = handleHTTP1ServerStream(b, tcpID, counterPair, superTimer, emitter, options, reqResMatcher) + switchingProtocolsHTTP2, err = handleHTTP1ServerStream(b, tcpID, counterPair, superTimer, emitter, options) if err == io.EOF || err == io.ErrUnexpectedEOF { break } else if err != nil { @@ -182,7 +181,7 @@ func (d dissecting) Dissect(b *bufio.Reader, isClient bool, tcpID *api.TcpID, co return nil } -func (d dissecting) Analyze(item *api.OutputChannelItem, resolvedSource string, resolvedDestination string, namespace string) *api.Entry { +func (d dissecting) Analyze(item *api.OutputChannelItem, resolvedSource string, resolvedDestination string) *api.Entry { var host, authority, path string request := item.Pair.Request.Payload.(map[string]interface{}) @@ -280,7 +279,6 @@ func (d dissecting) Analyze(item *api.OutputChannelItem, resolvedSource string, IP: item.ConnectionInfo.ServerIP, Port: item.ConnectionInfo.ServerPort, }, - Namespace: namespace, Outgoing: item.ConnectionInfo.IsOutgoing, Request: reqDetails, Response: resDetails, @@ -474,10 +472,6 @@ func (d dissecting) Macros() map[string]string { } } -func (d dissecting) NewResponseRequestMatcher() api.RequestResponseMatcher { - return createResponseRequestMatcher() -} - var Dissector dissecting func NewDissector() api.Dissector { diff --git a/tap/extensions/http/main_test.go b/tap/extensions/http/main_test.go index 90f90ab39..97cbc6430 100644 --- a/tap/extensions/http/main_test.go +++ b/tap/extensions/http/main_test.go @@ -11,6 +11,7 @@ import ( "os" "path" "path/filepath" + "sort" "testing" "time" @@ -38,6 +39,7 @@ func TestRegister(t *testing.T) { extension := &api.Extension{} dissector.Register(extension) assert.Equal(t, "http", extension.Protocol.Name) + assert.NotNil(t, extension.MatcherMap) } func TestMacros(t *testing.T) { @@ -121,8 +123,7 @@ func TestDissect(t *testing.T) { SrcPort: "1", DstPort: "2", } - reqResMatcher := dissector.NewResponseRequestMatcher() - err = dissector.Dissect(bufferClient, true, tcpIDClient, counterPair, &api.SuperTimer{}, superIdentifier, emitter, options, reqResMatcher) + err = dissector.Dissect(bufferClient, true, tcpIDClient, counterPair, &api.SuperTimer{}, superIdentifier, emitter, options) if err != nil && err != io.EOF && err != io.ErrUnexpectedEOF { panic(err) } @@ -140,7 +141,7 @@ func TestDissect(t *testing.T) { SrcPort: "2", DstPort: "1", } - err = dissector.Dissect(bufferServer, false, tcpIDServer, counterPair, &api.SuperTimer{}, superIdentifier, emitter, options, reqResMatcher) + err = dissector.Dissect(bufferServer, false, tcpIDServer, counterPair, &api.SuperTimer{}, superIdentifier, emitter, options) if err != nil && err != io.EOF && err != io.ErrUnexpectedEOF { panic(err) } @@ -154,6 +155,14 @@ func TestDissect(t *testing.T) { stop <- true + sort.Slice(items, func(i, j int) bool { + iMarshaled, err := json.Marshal(items[i]) + assert.Nil(t, err) + jMarshaled, err := json.Marshal(items[j]) + assert.Nil(t, err) + return len(iMarshaled) < len(jMarshaled) + }) + marshaled, err := json.Marshal(items) assert.Nil(t, err) @@ -205,7 +214,7 @@ func TestAnalyze(t *testing.T) { var entries []*api.Entry for _, item := range items { - entry := dissector.Analyze(item, "", "", "") + entry := dissector.Analyze(item, "", "") entries = append(entries, entry) } diff --git a/tap/extensions/http/matcher.go b/tap/extensions/http/matcher.go index 67c09136a..0a28a65ac 100644 --- a/tap/extensions/http/matcher.go +++ b/tap/extensions/http/matcher.go @@ -8,17 +8,16 @@ import ( "github.com/up9inc/mizu/tap/api" ) -// Key is {client_addr}_{client_port}_{dest_addr}_{dest_port}_{incremental_counter}_{proto_ident} +var reqResMatcher = createResponseRequestMatcher() // global + +// Key is {client_addr}:{client_port}->{dest_addr}:{dest_port}_{incremental_counter} type requestResponseMatcher struct { openMessagesMap *sync.Map } -func createResponseRequestMatcher() api.RequestResponseMatcher { - return &requestResponseMatcher{openMessagesMap: &sync.Map{}} -} - -func (matcher *requestResponseMatcher) GetMap() *sync.Map { - return matcher.openMessagesMap +func createResponseRequestMatcher() requestResponseMatcher { + newMatcher := &requestResponseMatcher{openMessagesMap: &sync.Map{}} + return *newMatcher } func (matcher *requestResponseMatcher) registerRequest(ident string, request *http.Request, captureTime time.Time, protoMinor int) *api.OutputChannelItem { diff --git a/tap/extensions/kafka/main.go b/tap/extensions/kafka/main.go index 0be1ccdbf..84a764064 100644 --- a/tap/extensions/kafka/main.go +++ b/tap/extensions/kafka/main.go @@ -33,27 +33,27 @@ type dissecting string func (d dissecting) Register(extension *api.Extension) { extension.Protocol = &_protocol + extension.MatcherMap = reqResMatcher.openMessagesMap } func (d dissecting) Ping() { log.Printf("pong %s", _protocol.Name) } -func (d dissecting) Dissect(b *bufio.Reader, isClient bool, tcpID *api.TcpID, counterPair *api.CounterPair, superTimer *api.SuperTimer, superIdentifier *api.SuperIdentifier, emitter api.Emitter, options *api.TrafficFilteringOptions, _reqResMatcher api.RequestResponseMatcher) error { - reqResMatcher := _reqResMatcher.(*requestResponseMatcher) +func (d dissecting) Dissect(b *bufio.Reader, isClient bool, tcpID *api.TcpID, counterPair *api.CounterPair, superTimer *api.SuperTimer, superIdentifier *api.SuperIdentifier, emitter api.Emitter, options *api.TrafficFilteringOptions) error { for { if superIdentifier.Protocol != nil && superIdentifier.Protocol != &_protocol { return errors.New("Identified by another protocol") } if isClient { - _, _, err := ReadRequest(b, tcpID, counterPair, superTimer, reqResMatcher) + _, _, err := ReadRequest(b, tcpID, counterPair, superTimer) if err != nil { return err } superIdentifier.Protocol = &_protocol } else { - err := ReadResponse(b, tcpID, counterPair, superTimer, emitter, reqResMatcher) + err := ReadResponse(b, tcpID, counterPair, superTimer, emitter) if err != nil { return err } @@ -62,7 +62,7 @@ func (d dissecting) Dissect(b *bufio.Reader, isClient bool, tcpID *api.TcpID, co } } -func (d dissecting) Analyze(item *api.OutputChannelItem, resolvedSource string, resolvedDestination string, namespace string) *api.Entry { +func (d dissecting) Analyze(item *api.OutputChannelItem, resolvedSource string, resolvedDestination string) *api.Entry { request := item.Pair.Request.Payload.(map[string]interface{}) reqDetails := request["details"].(map[string]interface{}) apiKey := ApiKey(reqDetails["apiKey"].(float64)) @@ -158,7 +158,6 @@ func (d dissecting) Analyze(item *api.OutputChannelItem, resolvedSource string, IP: item.ConnectionInfo.ServerIP, Port: item.ConnectionInfo.ServerPort, }, - Namespace: namespace, Outgoing: item.ConnectionInfo.IsOutgoing, Request: reqDetails, Response: item.Pair.Response.Payload.(map[string]interface{})["details"].(map[string]interface{}), @@ -216,10 +215,6 @@ func (d dissecting) Macros() map[string]string { } } -func (d dissecting) NewResponseRequestMatcher() api.RequestResponseMatcher { - return createResponseRequestMatcher() -} - var Dissector dissecting func NewDissector() api.Dissector { diff --git a/tap/extensions/kafka/matcher.go b/tap/extensions/kafka/matcher.go index d5c77618a..8bf8914bd 100644 --- a/tap/extensions/kafka/matcher.go +++ b/tap/extensions/kafka/matcher.go @@ -3,10 +3,9 @@ package kafka import ( "sync" "time" - - "github.com/up9inc/mizu/tap/api" ) +var reqResMatcher = CreateResponseRequestMatcher() // global const maxTry int = 3000 type RequestResponsePair struct { @@ -14,17 +13,14 @@ type RequestResponsePair struct { Response Response } -// Key is {client_addr}_{client_port}_{dest_addr}_{dest_port}_{correlation_id} +// Key is {client_addr}:{client_port}->{dest_addr}:{dest_port}::{correlation_id} type requestResponseMatcher struct { openMessagesMap *sync.Map } -func createResponseRequestMatcher() api.RequestResponseMatcher { - return &requestResponseMatcher{openMessagesMap: &sync.Map{}} -} - -func (matcher *requestResponseMatcher) GetMap() *sync.Map { - return matcher.openMessagesMap +func CreateResponseRequestMatcher() requestResponseMatcher { + newMatcher := &requestResponseMatcher{openMessagesMap: &sync.Map{}} + return *newMatcher } func (matcher *requestResponseMatcher) registerRequest(key string, request *Request) *RequestResponsePair { diff --git a/tap/extensions/kafka/request.go b/tap/extensions/kafka/request.go index 362d9e1df..982312936 100644 --- a/tap/extensions/kafka/request.go +++ b/tap/extensions/kafka/request.go @@ -19,7 +19,7 @@ type Request struct { CaptureTime time.Time `json:"captureTime"` } -func ReadRequest(r io.Reader, tcpID *api.TcpID, counterPair *api.CounterPair, superTimer *api.SuperTimer, reqResMatcher *requestResponseMatcher) (apiKey ApiKey, apiVersion int16, err error) { +func ReadRequest(r io.Reader, tcpID *api.TcpID, counterPair *api.CounterPair, superTimer *api.SuperTimer) (apiKey ApiKey, apiVersion int16, err error) { d := &decoder{reader: r, remain: 4} size := d.readInt32() @@ -214,7 +214,8 @@ func ReadRequest(r io.Reader, tcpID *api.TcpID, counterPair *api.CounterPair, su } key := fmt.Sprintf( - "%s_%s_%s_%s_%d", + "%d_%s:%s_%s:%s_%d", + counterPair.StreamId, tcpID.SrcIP, tcpID.SrcPort, tcpID.DstIP, diff --git a/tap/extensions/kafka/response.go b/tap/extensions/kafka/response.go index 0eb7950c7..dd9909034 100644 --- a/tap/extensions/kafka/response.go +++ b/tap/extensions/kafka/response.go @@ -16,7 +16,7 @@ type Response struct { CaptureTime time.Time `json:"captureTime"` } -func ReadResponse(r io.Reader, tcpID *api.TcpID, counterPair *api.CounterPair, superTimer *api.SuperTimer, emitter api.Emitter, reqResMatcher *requestResponseMatcher) (err error) { +func ReadResponse(r io.Reader, tcpID *api.TcpID, counterPair *api.CounterPair, superTimer *api.SuperTimer, emitter api.Emitter) (err error) { d := &decoder{reader: r, remain: 4} size := d.readInt32() @@ -44,7 +44,8 @@ func ReadResponse(r io.Reader, tcpID *api.TcpID, counterPair *api.CounterPair, s } key := fmt.Sprintf( - "%s_%s_%s_%s_%d", + "%d_%s:%s_%s:%s_%d", + counterPair.StreamId, tcpID.DstIP, tcpID.DstPort, tcpID.SrcIP, diff --git a/tap/extensions/redis/handlers.go b/tap/extensions/redis/handlers.go index c4eb7985e..a4a3a3858 100644 --- a/tap/extensions/redis/handlers.go +++ b/tap/extensions/redis/handlers.go @@ -6,14 +6,15 @@ import ( "github.com/up9inc/mizu/tap/api" ) -func handleClientStream(tcpID *api.TcpID, counterPair *api.CounterPair, superTimer *api.SuperTimer, emitter api.Emitter, request *RedisPacket, reqResMatcher *requestResponseMatcher) error { +func handleClientStream(tcpID *api.TcpID, counterPair *api.CounterPair, superTimer *api.SuperTimer, emitter api.Emitter, request *RedisPacket) error { counterPair.Lock() counterPair.Request++ requestCounter := counterPair.Request counterPair.Unlock() ident := fmt.Sprintf( - "%s_%s_%s_%s_%d", + "%d_%s:%s_%s:%s_%d", + counterPair.StreamId, tcpID.SrcIP, tcpID.DstIP, tcpID.SrcPort, @@ -35,14 +36,15 @@ func handleClientStream(tcpID *api.TcpID, counterPair *api.CounterPair, superTim return nil } -func handleServerStream(tcpID *api.TcpID, counterPair *api.CounterPair, superTimer *api.SuperTimer, emitter api.Emitter, response *RedisPacket, reqResMatcher *requestResponseMatcher) error { +func handleServerStream(tcpID *api.TcpID, counterPair *api.CounterPair, superTimer *api.SuperTimer, emitter api.Emitter, response *RedisPacket) error { counterPair.Lock() counterPair.Response++ responseCounter := counterPair.Response counterPair.Unlock() ident := fmt.Sprintf( - "%s_%s_%s_%s_%d", + "%d_%s:%s_%s:%s_%d", + counterPair.StreamId, tcpID.DstIP, tcpID.SrcIP, tcpID.DstPort, diff --git a/tap/extensions/redis/main.go b/tap/extensions/redis/main.go index db7480a7b..6de87739a 100644 --- a/tap/extensions/redis/main.go +++ b/tap/extensions/redis/main.go @@ -32,14 +32,14 @@ type dissecting string func (d dissecting) Register(extension *api.Extension) { extension.Protocol = &protocol + extension.MatcherMap = reqResMatcher.openMessagesMap } func (d dissecting) Ping() { log.Printf("pong %s", protocol.Name) } -func (d dissecting) Dissect(b *bufio.Reader, isClient bool, tcpID *api.TcpID, counterPair *api.CounterPair, superTimer *api.SuperTimer, superIdentifier *api.SuperIdentifier, emitter api.Emitter, options *api.TrafficFilteringOptions, _reqResMatcher api.RequestResponseMatcher) error { - reqResMatcher := _reqResMatcher.(*requestResponseMatcher) +func (d dissecting) Dissect(b *bufio.Reader, isClient bool, tcpID *api.TcpID, counterPair *api.CounterPair, superTimer *api.SuperTimer, superIdentifier *api.SuperIdentifier, emitter api.Emitter, options *api.TrafficFilteringOptions) error { is := &RedisInputStream{ Reader: b, Buf: make([]byte, 8192), @@ -52,9 +52,9 @@ func (d dissecting) Dissect(b *bufio.Reader, isClient bool, tcpID *api.TcpID, co } if isClient { - err = handleClientStream(tcpID, counterPair, superTimer, emitter, redisPacket, reqResMatcher) + err = handleClientStream(tcpID, counterPair, superTimer, emitter, redisPacket) } else { - err = handleServerStream(tcpID, counterPair, superTimer, emitter, redisPacket, reqResMatcher) + err = handleServerStream(tcpID, counterPair, superTimer, emitter, redisPacket) } if err != nil { @@ -63,7 +63,7 @@ func (d dissecting) Dissect(b *bufio.Reader, isClient bool, tcpID *api.TcpID, co } } -func (d dissecting) Analyze(item *api.OutputChannelItem, resolvedSource string, resolvedDestination string, namespace string) *api.Entry { +func (d dissecting) Analyze(item *api.OutputChannelItem, resolvedSource string, resolvedDestination string) *api.Entry { request := item.Pair.Request.Payload.(map[string]interface{}) response := item.Pair.Response.Payload.(map[string]interface{}) reqDetails := request["details"].(map[string]interface{}) @@ -96,7 +96,6 @@ func (d dissecting) Analyze(item *api.OutputChannelItem, resolvedSource string, IP: item.ConnectionInfo.ServerIP, Port: item.ConnectionInfo.ServerPort, }, - Namespace: namespace, Outgoing: item.ConnectionInfo.IsOutgoing, Request: reqDetails, Response: resDetails, @@ -128,10 +127,6 @@ func (d dissecting) Macros() map[string]string { } } -func (d dissecting) NewResponseRequestMatcher() api.RequestResponseMatcher { - return createResponseRequestMatcher() -} - var Dissector dissecting func NewDissector() api.Dissector { diff --git a/tap/extensions/redis/matcher.go b/tap/extensions/redis/matcher.go index e63c2f4b2..66e1c423b 100644 --- a/tap/extensions/redis/matcher.go +++ b/tap/extensions/redis/matcher.go @@ -7,17 +7,16 @@ import ( "github.com/up9inc/mizu/tap/api" ) -// Key is `{src_ip}_{dst_ip}_{src_ip}_{src_port}_{incremental_counter}` +var reqResMatcher = createResponseRequestMatcher() // global + +// Key is `{stream_id}_{src_ip}:{dst_ip}_{src_ip}:{src_port}_{incremental_counter}` type requestResponseMatcher struct { openMessagesMap *sync.Map } -func createResponseRequestMatcher() api.RequestResponseMatcher { - return &requestResponseMatcher{openMessagesMap: &sync.Map{}} -} - -func (matcher *requestResponseMatcher) GetMap() *sync.Map { - return matcher.openMessagesMap +func createResponseRequestMatcher() requestResponseMatcher { + newMatcher := &requestResponseMatcher{openMessagesMap: &sync.Map{}} + return *newMatcher } func (matcher *requestResponseMatcher) registerRequest(ident string, request *RedisPacket, captureTime time.Time) *api.OutputChannelItem { diff --git a/tap/passive_tapper.go b/tap/passive_tapper.go index 2779295e9..2a0a142ce 100644 --- a/tap/passive_tapper.go +++ b/tap/passive_tapper.go @@ -210,7 +210,6 @@ func startPassiveTapper(opts *TapOpts, outputItems chan *api.OutputChannelItem) assemblerMutex: &assembler.assemblerMutex, cleanPeriod: cleanPeriod, connectionTimeout: staleConnectionTimeout, - streamsMap: streamsMap, } cleaner.start() diff --git a/tap/tcp_reader.go b/tap/tcp_reader.go index ceb94e98e..bacda8ac8 100644 --- a/tap/tcp_reader.go +++ b/tap/tcp_reader.go @@ -47,7 +47,6 @@ type tcpReader struct { extension *api.Extension emitter api.Emitter counterPair *api.CounterPair - reqResMatcher api.RequestResponseMatcher sync.Mutex } @@ -95,7 +94,7 @@ func (h *tcpReader) Close() { func (h *tcpReader) run(wg *sync.WaitGroup) { defer wg.Done() b := bufio.NewReader(h) - err := h.extension.Dissector.Dissect(b, h.isClient, h.tcpID, h.counterPair, h.superTimer, h.parent.superIdentifier, h.emitter, filteringOptions, h.reqResMatcher) + err := h.extension.Dissector.Dissect(b, h.isClient, h.tcpID, h.counterPair, h.superTimer, h.parent.superIdentifier, h.emitter, filteringOptions) if err != nil { _, err = io.Copy(ioutil.Discard, b) if err != nil { diff --git a/tap/tcp_stream_factory.go b/tap/tcp_stream_factory.go index 9073e013b..ce56dcec8 100644 --- a/tap/tcp_stream_factory.go +++ b/tap/tcp_stream_factory.go @@ -29,9 +29,8 @@ type tcpStreamFactory struct { } type tcpStreamWrapper struct { - stream *tcpStream - reqResMatcher api.RequestResponseMatcher - createdAt time.Time + stream *tcpStream + createdAt time.Time } func NewTcpStreamFactory(emitter api.Emitter, streamsMap *tcpStreamMap, opts *TapOpts) *tcpStreamFactory { @@ -82,8 +81,8 @@ func (factory *tcpStreamFactory) New(net, transport gopacket.Flow, tcp *layers.T if stream.isTapTarget { stream.id = factory.streamsMap.nextId() for i, extension := range extensions { - reqResMatcher := extension.Dissector.NewResponseRequestMatcher() counterPair := &api.CounterPair{ + StreamId: stream.id, Request: 0, Response: 0, } @@ -104,7 +103,6 @@ func (factory *tcpStreamFactory) New(net, transport gopacket.Flow, tcp *layers.T extension: extension, emitter: factory.Emitter, counterPair: counterPair, - reqResMatcher: reqResMatcher, }) stream.servers = append(stream.servers, tcpReader{ msgQueue: make(chan tcpReaderDataMsg), @@ -123,13 +121,11 @@ func (factory *tcpStreamFactory) New(net, transport gopacket.Flow, tcp *layers.T extension: extension, emitter: factory.Emitter, counterPair: counterPair, - reqResMatcher: reqResMatcher, }) factory.streamsMap.Store(stream.id, &tcpStreamWrapper{ - stream: stream, - reqResMatcher: reqResMatcher, - createdAt: time.Now(), + stream: stream, + createdAt: time.Now(), }) factory.wg.Add(2) diff --git a/ui/src/components/Pages/TrafficPage/TrafficPage.tsx b/ui/src/components/Pages/TrafficPage/TrafficPage.tsx index d1f14985c..5b14a7d6c 100644 --- a/ui/src/components/Pages/TrafficPage/TrafficPage.tsx +++ b/ui/src/components/Pages/TrafficPage/TrafficPage.tsx @@ -76,7 +76,7 @@ export const TrafficPage: React.FC = ({setAnalyzeStatus}) => { const scrollableRef = useRef(null); const [openOasModal, setOpenOasModal] = useState(false); - + const handleOpenModal = () => setOpenOasModal(true); const handleCloseModal = () => setOpenOasModal(false); const [showTLSWarning, setShowTLSWarning] = useState(false); @@ -258,14 +258,8 @@ export const TrafficPage: React.FC = ({setAnalyzeStatus}) => { } } - const handleOpenOasModal = () => { - ws.current.close(); - setOpenOasModal(true); - } - const openServiceMapModalDebounce = debounce(() => { - ws.current.close(); - setServiceMapModalOpen(true); + setServiceMapModalOpen(true) }, 500); return ( @@ -291,7 +285,7 @@ export const TrafficPage: React.FC = ({setAnalyzeStatus}) => { variant="contained" className={commonClasses.outlinedButton + " " + commonClasses.imagedButton} style={{ marginRight: 25 }} - onClick={handleOpenOasModal} + onClick={handleOpenModal} > Show OAS }