diff --git a/core/core/scan.go b/core/core/scan.go index b893c709..2d835635 100644 --- a/core/core/scan.go +++ b/core/core/scan.go @@ -158,7 +158,7 @@ func (ks *Kubescape) Scan(ctx context.Context, scanInfo *cautils.ScanInfo) (*res // ===================== policies & resources ===================== ctxPolicies, spanPolicies := otel.Tracer("").Start(ctxInit, "policies & resources") policyHandler := policyhandler.NewPolicyHandler(interfaces.resourceHandler) - scanData, err := policyHandler.CollectResources(ctxPolicies, scanInfo.PolicyIdentifier, scanInfo) + scanData, err := policyHandler.CollectResources(ctxPolicies, scanInfo.PolicyIdentifier, scanInfo, cautils.NewProgressHandler("")) if err != nil { spanInit.End() return resultsHandling, err diff --git a/core/pkg/opaprocessor/processorhandler.go b/core/pkg/opaprocessor/processorhandler.go index 47ea2e27..6351553e 100644 --- a/core/pkg/opaprocessor/processorhandler.go +++ b/core/pkg/opaprocessor/processorhandler.go @@ -153,7 +153,7 @@ func (opap *OPAProcessor) Process(ctx context.Context, policies *cautils.Policie // processes rules for all controls in parallel for _, controlToPin := range policies.Controls { if progressListener != nil { - progressListener.ProgressJob(1, fmt.Sprintf("Control %s", controlToPin.ControlID)) + progressListener.ProgressJob(1, fmt.Sprintf("Control: %s", controlToPin.ControlID)) } control := controlToPin diff --git a/core/pkg/policyhandler/handlenotification.go b/core/pkg/policyhandler/handlenotification.go index 7058e4ac..14afb35b 100644 --- a/core/pkg/policyhandler/handlenotification.go +++ b/core/pkg/policyhandler/handlenotification.go @@ -12,6 +12,7 @@ import ( clientcmdapi "k8s.io/client-go/tools/clientcmd/api" cloudsupportv1 "github.com/kubescape/k8s-interface/cloudsupport/v1" + "github.com/kubescape/kubescape/v2/core/pkg/opaprocessor" reportv2 "github.com/kubescape/opa-utils/reporthandling/v2" "github.com/kubescape/k8s-interface/cloudsupport" @@ -35,7 +36,7 @@ func NewPolicyHandler(resourceHandler resourcehandler.IResourceHandler) *PolicyH } } -func (policyHandler *PolicyHandler) CollectResources(ctx context.Context, policyIdentifier []cautils.PolicyIdentifier, scanInfo *cautils.ScanInfo) (*cautils.OPASessionObj, error) { +func (policyHandler *PolicyHandler) CollectResources(ctx context.Context, policyIdentifier []cautils.PolicyIdentifier, scanInfo *cautils.ScanInfo, progressListener opaprocessor.IJobProgressNotificationClient) (*cautils.OPASessionObj, error) { opaSessionObj := cautils.NewOPASessionObj(ctx, nil, nil, scanInfo) // validate notification @@ -47,7 +48,7 @@ func (policyHandler *PolicyHandler) CollectResources(ctx context.Context, policy return opaSessionObj, err } - err := policyHandler.getResources(ctx, policyIdentifier, opaSessionObj) + err := policyHandler.getResources(ctx, policyIdentifier, opaSessionObj, progressListener) if err != nil { return opaSessionObj, err } @@ -59,7 +60,7 @@ func (policyHandler *PolicyHandler) CollectResources(ctx context.Context, policy return opaSessionObj, nil } -func (policyHandler *PolicyHandler) getResources(ctx context.Context, policyIdentifier []cautils.PolicyIdentifier, opaSessionObj *cautils.OPASessionObj) error { +func (policyHandler *PolicyHandler) getResources(ctx context.Context, policyIdentifier []cautils.PolicyIdentifier, opaSessionObj *cautils.OPASessionObj, progressListener opaprocessor.IJobProgressNotificationClient) error { ctx, span := otel.Tracer("").Start(ctx, "policyHandler.getResources") defer span.End() opaSessionObj.Report.ClusterAPIServerInfo = policyHandler.resourceHandler.GetClusterAPIServerInfo(ctx) @@ -69,7 +70,7 @@ func (policyHandler *PolicyHandler) getResources(ctx context.Context, policyIden setCloudMetadata(opaSessionObj) } - resourcesMap, allResources, ksResources, err := policyHandler.resourceHandler.GetResources(ctx, opaSessionObj, &policyIdentifier[0].Designators) + resourcesMap, allResources, ksResources, err := policyHandler.resourceHandler.GetResources(ctx, opaSessionObj, &policyIdentifier[0].Designators, progressListener) if err != nil { return err } diff --git a/core/pkg/policyhandler/handlenotification_test.go b/core/pkg/policyhandler/handlenotification_test.go index ba772d1e..57a96436 100644 --- a/core/pkg/policyhandler/handlenotification_test.go +++ b/core/pkg/policyhandler/handlenotification_test.go @@ -249,12 +249,12 @@ func Test_getResources(t *testing.T) { policyIdentifier := []cautils.PolicyIdentifier{{}} assert.NotPanics(t, func() { - policyHandler.getResources(context.TODO(), policyIdentifier, objSession) + policyHandler.getResources(context.TODO(), policyIdentifier, objSession, cautils.NewProgressHandler("")) }, "Cluster named .*eks.* without a cloud config panics on cluster scan !") assert.NotPanics(t, func() { objSession.Metadata.ScanMetadata.ScanningTarget = reportv2.File - policyHandler.getResources(context.TODO(), policyIdentifier, objSession) + policyHandler.getResources(context.TODO(), policyIdentifier, objSession, cautils.NewProgressHandler("")) }, "Cluster named .*eks.* without a cloud config panics on non-cluster scan !") } diff --git a/core/pkg/resourcehandler/filesloader.go b/core/pkg/resourcehandler/filesloader.go index e823dd55..cb525c62 100644 --- a/core/pkg/resourcehandler/filesloader.go +++ b/core/pkg/resourcehandler/filesloader.go @@ -15,6 +15,7 @@ import ( "github.com/kubescape/go-logger/helpers" "github.com/kubescape/k8s-interface/k8sinterface" "github.com/kubescape/kubescape/v2/core/cautils" + "github.com/kubescape/kubescape/v2/core/pkg/opaprocessor" ) // FileResourceHandler handle resources from files and URLs @@ -31,7 +32,7 @@ func NewFileResourceHandler(_ context.Context, inputPatterns []string, registryA } } -func (fileHandler *FileResourceHandler) GetResources(ctx context.Context, sessionObj *cautils.OPASessionObj, _ *armotypes.PortalDesignator) (*cautils.K8SResources, map[string]workloadinterface.IMetadata, *cautils.KSResources, error) { +func (fileHandler *FileResourceHandler) GetResources(ctx context.Context, sessionObj *cautils.OPASessionObj, _ *armotypes.PortalDesignator, progressListener opaprocessor.IJobProgressNotificationClient) (*cautils.K8SResources, map[string]workloadinterface.IMetadata, *cautils.KSResources, error) { // // build resources map diff --git a/core/pkg/resourcehandler/k8sresources.go b/core/pkg/resourcehandler/k8sresources.go index 6865181e..d05eb0a4 100644 --- a/core/pkg/resourcehandler/k8sresources.go +++ b/core/pkg/resourcehandler/k8sresources.go @@ -9,6 +9,7 @@ import ( "github.com/kubescape/go-logger/helpers" "github.com/kubescape/kubescape/v2/core/cautils" "github.com/kubescape/kubescape/v2/core/pkg/hostsensorutils" + "github.com/kubescape/kubescape/v2/core/pkg/opaprocessor" "github.com/kubescape/opa-utils/objectsenvelopes" "github.com/kubescape/opa-utils/reporthandling/apis" @@ -56,7 +57,7 @@ func NewK8sResourceHandler(k8s *k8sinterface.KubernetesApi, fieldSelector IField } } -func (k8sHandler *K8sResourceHandler) GetResources(ctx context.Context, sessionObj *cautils.OPASessionObj, designator *armotypes.PortalDesignator) (*cautils.K8SResources, map[string]workloadinterface.IMetadata, *cautils.KSResources, error) { +func (k8sHandler *K8sResourceHandler) GetResources(ctx context.Context, sessionObj *cautils.OPASessionObj, designator *armotypes.PortalDesignator, progressListener opaprocessor.IJobProgressNotificationClient) (*cautils.K8SResources, map[string]workloadinterface.IMetadata, *cautils.KSResources, error) { allResources := map[string]workloadinterface.IMetadata{} // get k8s resources @@ -143,7 +144,7 @@ func (k8sHandler *K8sResourceHandler) GetResources(ctx context.Context, sessionO // check that controls use cloud resources if len(cloudResources) > 0 { - err := k8sHandler.collectCloudResources(ctx, sessionObj, allResources, ksResourceMap, cloudResources) + err := k8sHandler.collectCloudResources(ctx, sessionObj, allResources, ksResourceMap, cloudResources, progressListener) if err != nil { cautils.SetInfoMapForResources(err.Error(), cloudResources, sessionObj.InfoMap) logger.L().Debug("failed to collect cloud data", helpers.Error(err)) @@ -153,7 +154,7 @@ func (k8sHandler *K8sResourceHandler) GetResources(ctx context.Context, sessionO return k8sResourcesMap, allResources, ksResourceMap, nil } -func (k8sHandler *K8sResourceHandler) collectCloudResources(ctx context.Context, sessionObj *cautils.OPASessionObj, allResources map[string]workloadinterface.IMetadata, ksResourceMap *cautils.KSResources, cloudResources []string) error { +func (k8sHandler *K8sResourceHandler) collectCloudResources(ctx context.Context, sessionObj *cautils.OPASessionObj, allResources map[string]workloadinterface.IMetadata, ksResourceMap *cautils.KSResources, cloudResources []string, progressListener opaprocessor.IJobProgressNotificationClient) error { clusterName := cautils.ClusterName provider := cloudsupport.GetCloudProvider(clusterName) if provider == "" { @@ -165,7 +166,17 @@ func (k8sHandler *K8sResourceHandler) collectCloudResources(ctx context.Context, } logger.L().Debug("cloud", helpers.String("cluster", clusterName), helpers.String("clusterName", clusterName), helpers.String("provider", provider)) + logger.L().Info("Downloading cloud resources") + // start progressbar during pull of cloud resources (this can take a while). + if progressListener != nil { + progressListener.Start(len(cloudResources)) + defer progressListener.Stop() + } for resourceKind, resourceGetter := range cloudResourceGetterMapping { + // set way to progress + if progressListener != nil { + progressListener.ProgressJob(1, fmt.Sprintf("Cloud Resource: %s", resourceKind)) + } if !cloudResourceRequired(cloudResources, resourceKind) { continue } @@ -186,6 +197,7 @@ func (k8sHandler *K8sResourceHandler) collectCloudResources(ctx context.Context, allResources[wl.GetID()] = wl (*ksResourceMap)[fmt.Sprintf("%s/%s", wl.GetApiVersion(), wl.GetKind())] = []string{wl.GetID()} } + logger.L().Success("Downloaded cloud resources") // get api server info resource if cloudResourceRequired(cloudResources, string(cloudsupport.TypeApiServerInfo)) { diff --git a/core/pkg/resourcehandler/resourceshandler.go b/core/pkg/resourcehandler/resourceshandler.go index 7ed2de35..651aa197 100644 --- a/core/pkg/resourcehandler/resourceshandler.go +++ b/core/pkg/resourcehandler/resourceshandler.go @@ -6,10 +6,11 @@ import ( "github.com/armosec/armoapi-go/armotypes" "github.com/kubescape/k8s-interface/workloadinterface" "github.com/kubescape/kubescape/v2/core/cautils" + "github.com/kubescape/kubescape/v2/core/pkg/opaprocessor" "k8s.io/apimachinery/pkg/version" ) type IResourceHandler interface { - GetResources(context.Context, *cautils.OPASessionObj, *armotypes.PortalDesignator) (*cautils.K8SResources, map[string]workloadinterface.IMetadata, *cautils.KSResources, error) + GetResources(context.Context, *cautils.OPASessionObj, *armotypes.PortalDesignator, opaprocessor.IJobProgressNotificationClient) (*cautils.K8SResources, map[string]workloadinterface.IMetadata, *cautils.KSResources, error) GetClusterAPIServerInfo(ctx context.Context) *version.Info }