k8scommon_impl

package
v0.59.0 Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Sep 17, 2026 License: Apache-2.0 Imports: 20 Imported by: 0

Documentation

Index

Constants

This section is empty.

Variables

View Source
var AutocompleteClusterIdentityTask = inspectiontaskbase.NewGlobalCachedTask(k8scommon.AutocompleteClusterIdentityTaskID, []coretask.Dependency{
	k8scommon.ClusterNamePrefixTaskRef,
	gcpcommon.InputProjectIdTaskID.Ref(),
	gcpcommon.InputStartTimeTaskID.Ref(),
	gcpcommon.InputEndTimeTaskID.Ref(),
	k8scommon.AutocompleteMetricsK8sContainerTaskID.Ref(),
	gcpcommon.APIClientFactoryTaskID.Ref(),
	gcpcommon.APIClientCallOptionsInjectorTaskID.Ref(),
}, func(ctx context.Context, prevValue inspectiontaskbase.CacheableTaskResult[*inspectioncore.AutocompleteResult[k8scommon.GoogleCloudClusterIdentity]]) (inspectiontaskbase.CacheableTaskResult[*inspectioncore.AutocompleteResult[k8scommon.GoogleCloudClusterIdentity]], error) {
	clusterNamePrefix := coretask.GetTaskResult(ctx, k8scommon.ClusterNamePrefixTaskRef)
	projectID := coretask.GetTaskResult(ctx, gcpcommon.InputProjectIdTaskID.Ref())
	startTime := coretask.GetTaskResult(ctx, gcpcommon.InputStartTimeTaskID.Ref())
	endTime := coretask.GetTaskResult(ctx, gcpcommon.InputEndTimeTaskID.Ref())
	metricsType := coretask.GetTaskResult(ctx, k8scommon.AutocompleteMetricsK8sContainerTaskID.Ref())
	cf := coretask.GetTaskResult(ctx, gcpcommon.APIClientFactoryTaskID.Ref())
	optionInjector := coretask.GetTaskResult(ctx, gcpcommon.APIClientCallOptionsInjectorTaskID.Ref())

	currentDigest := fmt.Sprintf("%s-%s-%d-%d", clusterNamePrefix.PrefixFor(k8scommon.ClusterNameUsageK8sCluster), projectID, startTime.Unix(), endTime.Unix())
	if currentDigest == prevValue.DependencyDigest {
		return prevValue, nil
	}
	if projectID == "" {
		return inspectiontaskbase.CacheableTaskResult[*inspectioncore.AutocompleteResult[k8scommon.GoogleCloudClusterIdentity]]{
			Value: &inspectioncore.AutocompleteResult[k8scommon.GoogleCloudClusterIdentity]{
				Values: []k8scommon.GoogleCloudClusterIdentity{},
				Error:  "",
				Hint:   "Cluster names are suggested after the project ID is provided.",
			},
			DependencyDigest: currentDigest,
		}, nil
	}

	errorString := ""
	hintString := ""
	if endTime.Before(time.Now().Add(-time.Hour * 24 * 30 * 24)) {
		hintString = "The end time is more than 24 months ago. Suggested cluster names may not be complete."
	}

	client, err := cf.MonitoringMetricClient(ctx, googlecloud.Project(projectID))
	if err != nil {
		return prevValue, fmt.Errorf("failed to create monitoring metric client: %w", err)
	}

	ctx = optionInjector.InjectToCallContext(ctx, googlecloud.Project(projectID))
	filter := fmt.Sprintf(`metric.type="%s" AND resource.type="k8s_container"`, metricsType)
	metricsLabels, err := googlecloud.QueryResourceLabelsFromMetrics(ctx, client, projectID, filter, startTime, endTime, []string{"resource.label.cluster_name", "resource.label.location"})
	if err != nil {
		errorString = err.Error()
	}
	metricsLabels = filterAndTrimPrefixFromClusterNames(metricsLabels, clusterNamePrefix.PrefixFor(k8scommon.ClusterNameUsageK8sCluster))
	if hintString == "" && errorString == "" && len(metricsLabels) == 0 {
		hintString = fmt.Sprintf("No cluster names found between %s and %s. It is highly likely that the time range is incorrect. Please verify the time range, or proceed by manually entering the cluster name.", startTime.Format(time.RFC3339), endTime.Format(time.RFC3339))
	}

	identities := make([]k8scommon.GoogleCloudClusterIdentity, len(metricsLabels))
	for i, labels := range metricsLabels {
		identities[i] = k8scommon.GoogleCloudClusterIdentity{
			ProjectID:    projectID,
			PrefixPolicy: clusterNamePrefix,
			ClusterName:  labels["cluster_name"],
			Location:     labels["location"],
		}
	}

	return inspectiontaskbase.CacheableTaskResult[*inspectioncore.AutocompleteResult[k8scommon.GoogleCloudClusterIdentity]]{
		DependencyDigest: currentDigest,
		Value: &inspectioncore.AutocompleteResult[k8scommon.GoogleCloudClusterIdentity]{
			Values: identities,
			Error:  errorString,
			Hint:   hintString,
		},
	}, nil
})
View Source
var AutocompleteLocationForClusterTask = inspectiontaskbase.NewGlobalCachedTask(k8scommon.AutocompleteLocationForClusterTaskID, []coretask.Dependency{
	k8scommon.InputClusterNameTaskID.Ref(),
	gcpcommon.InputProjectIdTaskID.Ref(),
	gcpcommon.InputStartTimeTaskID.Ref(),
	gcpcommon.InputEndTimeTaskID.Ref(),
	k8scommon.AutocompleteClusterIdentityTaskID.Ref(),
}, func(ctx context.Context, prevValue inspectiontaskbase.CacheableTaskResult[*inspectioncore.AutocompleteResult[string]]) (inspectiontaskbase.CacheableTaskResult[*inspectioncore.AutocompleteResult[string]], error) {
	projectID := coretask.GetTaskResult(ctx, gcpcommon.InputProjectIdTaskID.Ref())
	clusterName := coretask.GetTaskResult(ctx, k8scommon.InputClusterNameTaskID.Ref())
	startTime := coretask.GetTaskResult(ctx, gcpcommon.InputStartTimeTaskID.Ref())
	endTime := coretask.GetTaskResult(ctx, gcpcommon.InputEndTimeTaskID.Ref())
	clusterIdentities := coretask.GetTaskResult(ctx, k8scommon.AutocompleteClusterIdentityTaskID.Ref())

	currentDigest := fmt.Sprintf("%s-%s-%d-%d", clusterName, projectID, startTime.Unix(), endTime.Unix())
	if currentDigest == prevValue.DependencyDigest {
		return prevValue, nil
	}
	if projectID == "" {
		return inspectiontaskbase.CacheableTaskResult[*inspectioncore.AutocompleteResult[string]]{
			Value: &inspectioncore.AutocompleteResult[string]{
				Values: []string{},
				Error:  "",
				Hint:   "Locations will be suggested after the project ID is provided.",
			},
			DependencyDigest: currentDigest,
		}, nil
	}
	if clusterIdentities.Error != "" {
		return inspectiontaskbase.CacheableTaskResult[*inspectioncore.AutocompleteResult[string]]{
			Value: &inspectioncore.AutocompleteResult[string]{
				Values: []string{},
				Error:  clusterIdentities.Error,
				Hint:   clusterIdentities.Hint,
			},
			DependencyDigest: currentDigest,
		}, nil
	}
	if clusterName == "" {
		return inspectiontaskbase.CacheableTaskResult[*inspectioncore.AutocompleteResult[string]]{
			Value: &inspectioncore.AutocompleteResult[string]{
				Values: []string{},
				Error:  "",
				Hint:   "Locations will be suggested after the cluster name is provided.",
			},
			DependencyDigest: currentDigest,
		}, nil
	}
	result := &inspectioncore.AutocompleteResult[string]{
		Values: []string{},
		Error:  "",
		Hint:   "",
	}

	for _, identity := range clusterIdentities.Values {
		if identity.ClusterName == clusterName {
			result.Values = append(result.Values, identity.Location)
		}
	}
	return inspectiontaskbase.CacheableTaskResult[*inspectioncore.AutocompleteResult[string]]{
		Value:            result,
		DependencyDigest: currentDigest,
	}, nil
}, coretask.WithSelectionPriority(500))

AutocompleteLocationForClusterTask returns the location for the given cluster name.

View Source
var AutocompleteMetricsK8sContainerTask = coretask.NewTask(k8scommon.AutocompleteMetricsK8sContainerTaskID, []coretask.Dependency{}, func(ctx context.Context) (string, error) {

	return "kubernetes.io/anthos/up", nil
})

AutocompleteMetricsK8sContainerTask is the task to provide the default metrics type to collect the cluster names. The resource type "k8s_container" must be available on the returned metrics type. This task is overridden in GKE clusters.

View Source
var AutocompleteMetricsK8sNodeTask = coretask.NewTask(k8scommon.AutocompleteMetricsK8sNodeTaskID, []coretask.Dependency{}, func(ctx context.Context) (string, error) {
	return "kubernetes.io/anthos/up", nil
})
View Source
var AutocompleteNamespacesTask = inspectiontaskbase.NewGlobalCachedTask(k8scommon.AutocompleteNamespacesTaskID, []coretask.Dependency{
	k8scommon.ClusterIdentityTaskID.Ref(),
	gcpcommon.InputStartTimeTaskID.Ref(),
	gcpcommon.InputEndTimeTaskID.Ref(),
	gcpcommon.APIClientFactoryTaskID.Ref(),
	gcpcommon.APIClientCallOptionsInjectorTaskID.Ref(),
	k8scommon.AutocompleteMetricsK8sContainerTaskID.Ref(),
}, func(ctx context.Context, prevValue inspectiontaskbase.CacheableTaskResult[*inspectioncore.AutocompleteResult[string]]) (inspectiontaskbase.CacheableTaskResult[*inspectioncore.AutocompleteResult[string]], error) {
	metricsType := coretask.GetTaskResult(ctx, k8scommon.AutocompleteMetricsK8sContainerTaskID.Ref())
	return queryClusterScopedAutocompleteMetrics(ctx, prevValue, metricsType, clusterScopedAutocompleteConfig{
		resourceType:       "k8s_container",
		resourceLabelKey:   "namespace_name",
		targetNameSingular: "namespace name",
		targetNamePlural:   "namespace names",
	})
})
View Source
var AutocompleteNodeNamesTask = inspectiontaskbase.NewGlobalCachedTask(k8scommon.AutocompleteNodeNamesTaskID, []coretask.Dependency{
	k8scommon.ClusterIdentityTaskID.Ref(),
	gcpcommon.InputStartTimeTaskID.Ref(),
	gcpcommon.InputEndTimeTaskID.Ref(),
	gcpcommon.APIClientFactoryTaskID.Ref(),
	gcpcommon.APIClientCallOptionsInjectorTaskID.Ref(),
	k8scommon.AutocompleteMetricsK8sNodeTaskID.Ref(),
}, func(ctx context.Context, prevValue inspectiontaskbase.CacheableTaskResult[*inspectioncore.AutocompleteResult[string]]) (inspectiontaskbase.CacheableTaskResult[*inspectioncore.AutocompleteResult[string]], error) {
	metricsType := coretask.GetTaskResult(ctx, k8scommon.AutocompleteMetricsK8sNodeTaskID.Ref())
	return queryClusterScopedAutocompleteMetrics(ctx, prevValue, metricsType, clusterScopedAutocompleteConfig{
		resourceType:       "k8s_node",
		resourceLabelKey:   "node_name",
		targetNameSingular: "node name",
		targetNamePlural:   "node names",
	})
})
View Source
var AutocompletePodNamesTask = inspectiontaskbase.NewGlobalCachedTask(k8scommon.AutocompletePodNamesTaskID, []coretask.Dependency{
	k8scommon.ClusterIdentityTaskID.Ref(),
	gcpcommon.InputStartTimeTaskID.Ref(),
	gcpcommon.InputEndTimeTaskID.Ref(),
	gcpcommon.APIClientFactoryTaskID.Ref(),
	gcpcommon.APIClientCallOptionsInjectorTaskID.Ref(),
	k8scommon.AutocompleteMetricsK8sContainerTaskID.Ref(),
}, func(ctx context.Context, prevValue inspectiontaskbase.CacheableTaskResult[*inspectioncore.AutocompleteResult[string]]) (inspectiontaskbase.CacheableTaskResult[*inspectioncore.AutocompleteResult[string]], error) {
	metricsType := coretask.GetTaskResult(ctx, k8scommon.AutocompleteMetricsK8sContainerTaskID.Ref())
	return queryClusterScopedAutocompleteMetrics(ctx, prevValue, metricsType, clusterScopedAutocompleteConfig{
		resourceType:       "k8s_container",
		resourceLabelKey:   "pod_name",
		targetNameSingular: "pod name",
		targetNamePlural:   "pod names",
	})
})

HeaderSuggestedFileNameTask is a task to supply the suggested file name of the KHI file generated. This name is used in frontend to save the inspection data as a file.

View Source
var InputClusterNameTask = formtask.NewTextFormTaskBuilder(k8scommon.InputClusterNameTaskID, gcpcommon.PriorityForResourceIdentifierGroup+4000, "Cluster name").
	WithDependencies([]coretask.Dependency{k8scommon.AutocompleteClusterIdentityTaskID.Ref()}).
	WithDescription("The cluster name to gather logs.").
	WithDefaultValueFunc(func(ctx context.Context, previousValues []string) (string, error) {
		clusters := coretask.GetTaskResult(ctx, k8scommon.AutocompleteClusterIdentityTaskID.Ref())

		if len(previousValues) > 0 && hasClusterNameInAutocomplete(clusters.Values, previousValues[0]) {
			return previousValues[0], nil
		}
		if len(clusters.Values) == 0 {
			return "", nil
		}
		return clusters.Values[0].ClusterName, nil
	}).
	WithSuggestionsFunc(func(ctx context.Context, value string, previousValues []string) ([]string, error) {
		clusters := coretask.GetTaskResult(ctx, k8scommon.AutocompleteClusterIdentityTaskID.Ref())
		return common.SortForAutocomplete(value, dedupeClusterName(clusters.Values)), nil
	}).
	WithHintFunc(func(ctx context.Context, value string, convertedValue any) (string, inspectionmetadata.ParameterHintType, error) {
		clusters := coretask.GetTaskResult(ctx, k8scommon.AutocompleteClusterIdentityTaskID.Ref())

		if clusters.Error != "" {
			return fmt.Sprintf("Failed to obtain the cluster list due to the error '%s'.\n The suggestion list won't popup", clusters.Error), inspectionmetadata.Warning, nil
		}
		if clusters.Hint != "" {
			return clusters.Hint, inspectionmetadata.Info, nil
		}
		for _, suggestedCluster := range clusters.Values {
			if suggestedCluster.ClusterName == convertedValue.(string) {
				return "", inspectionmetadata.Info, nil
			}
		}

		availableClusterNameStr := ""
		for _, cluster := range dedupeClusterName(clusters.Values) {
			availableClusterNameStr += fmt.Sprintf("* %s\n", cluster)
		}
		return fmt.Sprintf("Cluster '%s' was not found in the specified project at this time. It works for the clusters existed in the past but make sure the cluster name is right if you believe the cluster should be there.\nAvailable cluster names:\n%s", value, availableClusterNameStr), inspectionmetadata.Warning, nil
	}).
	WithValidator(validateClusterName).
	WithConverter(convertClusterName).
	Build()

InputClusterNameTask is a form task receiving cluster name from the user. This task returns the raw cluster name without the prefix defined from the cluster type. This input also supports autocomplete cluster names from some task having ID for k8scommon.AutocompleteClusterIdentityTaskID.

View Source
var InputKindFilterTask = formtask.NewSetFormTaskBuilder(k8scommon.InputKindFilterTaskID, gcpcommon.PriorityForK8sResourceFilterGroup+5000, "Kind").
	WithDefaultValueConstant([]string{"@any", "-leases"}, true).
	WithDescription("The kinds of resources to gather logs. Specify `@any` to query all kinds of resources, or prefix with `-` to exclude specific kinds (e.g., `-leases`). `@legacy_default` matches a set of kinds frequently queried in legacy KHI versions.").
	WithAllowAddAll(false).
	WithAllowRemoveAll(false).
	WithAllowCustomValue(true).
	WithOptionsFunc(func(ctx context.Context, previousValues []string) ([]inspectionmetadata.SetParameterFormFieldOptionItem, error) {
		return []inspectionmetadata.SetParameterFormFieldOptionItem{
			{ID: "@any", Description: "[Alias] An alias matches any of the kinds"},
			{ID: "@legacy_default", Description: "[Alias] An alias matches a set of kinds frequently queried in legacy KHI versions."},
		}, nil
	}).
	WithValidator(func(ctx context.Context, value []string) (string, error) {
		if len(value) == 0 {
			return "kind filter can't be empty", nil
		}
		filterInStr := strings.Join(value, " ")
		result, err := gcpqueryutil.ParseSetFilter(filterInStr, inputKindsAliasMap, true, true, true)
		if err != nil {
			return "", err
		}
		return result.ValidationError, nil
	}).
	WithConverter(func(ctx context.Context, value []string) (*gcpqueryutil.SetFilterParseResult, error) {
		filterInStr := strings.Join(value, " ")
		result, err := gcpqueryutil.ParseSetFilter(filterInStr, inputKindsAliasMap, true, true, true)
		if err != nil {
			return nil, err
		}
		return result, nil
	}).
	Build()

InputKindFilterTask is a form task for inputting the kind filter.

View Source
var InputNamespaceFilterTask = formtask.NewSetFormTaskBuilder(k8scommon.InputNamespaceFilterTaskID, gcpcommon.PriorityForK8sResourceFilterGroup+4000, "Namespaces").
	WithDependencies([]coretask.Dependency{k8scommon.AutocompleteNamespacesTaskID.Ref()}).
	WithDefaultValueConstant([]string{"@all_cluster_scoped", "@all_namespaced"}, true).
	WithDescription("The namespace of resources to gather logs. Specify `@all_cluster_scoped` to gather logs for all non-namespaced resources. Specify `@all_namespaced` to gather logs for all namespaced resources.").
	WithAllowAddAll(false).
	WithAllowRemoveAll(false).
	WithAllowCustomValue(true).
	WithOptionsFunc(func(ctx context.Context, previousValues []string) ([]inspectionmetadata.SetParameterFormFieldOptionItem, error) {
		result := []inspectionmetadata.SetParameterFormFieldOptionItem{
			{ID: "@all_cluster_scoped", Description: "[Alias] An alias matches any of the cluster scoped resources"},
			{ID: "@all_namespaced", Description: "[Alias] An alias matches any of the namespaced resources"},
		}
		namespaces := coretask.GetTaskResult(ctx, k8scommon.AutocompleteNamespacesTaskID.Ref())
		for index, namespace := range namespaces.Values {
			if index >= maxNamespaceFilterOptions {
				break
			}
			result = append(result, inspectionmetadata.SetParameterFormFieldOptionItem{ID: namespace})
		}
		return result, nil
	}).
	WithHintFunc(func(ctx context.Context, value []string, convertedValue any) (string, inspectionmetadata.ParameterHintType, error) {
		namespaces := coretask.GetTaskResult(ctx, k8scommon.AutocompleteNamespacesTaskID.Ref())
		if len(namespaces.Values) > maxNamespaceFilterOptions {
			return fmt.Sprintf("Some namespaces are not shown on the suggestion list because the number of namespaces is %d, which is more than %d.", len(namespaces.Values), maxNamespaceFilterOptions), inspectionmetadata.Warning, nil
		}
		return "", inspectionmetadata.None, nil
	}).
	WithValidator(func(ctx context.Context, value []string) (string, error) {
		if len(value) == 0 {
			return "namespace filter can't be empty", nil
		}
		namespaceFilterInStr := strings.Join(value, " ")
		result, err := gcpqueryutil.ParseSetFilter(namespaceFilterInStr, inputNamespacesAliasMap, false, false, true)
		if err != nil {
			return "", err
		}
		return result.ValidationError, nil
	}).
	WithConverter(func(ctx context.Context, value []string) (*gcpqueryutil.SetFilterParseResult, error) {
		namespaceFilterInStr := strings.Join(value, " ")
		result, err := gcpqueryutil.ParseSetFilter(namespaceFilterInStr, inputNamespacesAliasMap, false, false, true)
		if err != nil {
			return nil, err
		}
		return result, nil
	}).
	Build()

InputNamespaceFilterTask is a form task for inputting the namespace filter.

View Source
var InputNodeNameFilterTask = formtask.NewSetFormTaskBuilder(k8scommon.InputNodeNameFilterTaskID, gcpcommon.PriorityForK8sResourceFilterGroup+3000, "Node names").
	WithDependencies([]coretask.Dependency{k8scommon.AutocompleteNodeNamesTaskID.Ref()}).
	WithDefaultValueConstant([]string{}, true).
	WithDescription("A space-separated list of node name substrings used to collect node-related logs. If left blank, KHI gathers logs from all nodes in the cluster.").
	WithAllowAddAll(false).
	WithAllowRemoveAll(false).
	WithAllowCustomValue(true).
	WithValidator(func(ctx context.Context, value []string) (string, error) {
		for _, v := range value {
			if !nodeNameSubstringValidator.MatchString(v) {
				return fmt.Sprintf("invalid node name substring: %s", v), nil
			}
		}
		return "", nil
	}).
	WithOptionsFunc(func(ctx context.Context, prevValue []string) ([]inspectionmetadata.SetParameterFormFieldOptionItem, error) {
		result := []inspectionmetadata.SetParameterFormFieldOptionItem{}
		nodeNames := coretask.GetTaskResult(ctx, k8scommon.AutocompleteNodeNamesTaskID.Ref())
		for i, v := range nodeNames.Values {
			if i >= maxNodeNameFilterOptions {
				break
			}
			result = append(result, inspectionmetadata.SetParameterFormFieldOptionItem{ID: v})
		}
		return result, nil
	}).
	WithHintFunc(func(ctx context.Context, value []string, convertedValue any) (string, inspectionmetadata.ParameterHintType, error) {
		nodeNames := coretask.GetTaskResult(ctx, k8scommon.AutocompleteNodeNamesTaskID.Ref())
		if len(nodeNames.Values) > maxNodeNameFilterOptions {
			return fmt.Sprintf("Some node names are not shown on the suggestion list because the number of node names is %d, which is more than %d.", len(nodeNames.Values), maxNodeNameFilterOptions), inspectionmetadata.Warning, nil
		}
		return "", inspectionmetadata.None, nil
	}).
	Build()

InputNodeNameFilterTask is a task to collect list of substrings of node names. This input value is used in querying k8s_node or serialport logs.

View Source
var NEGNamesDiscoveryTask = inspectiontaskbase.NewInspectionTask(
	k8scommon.NEGNamesDiscoveryTaskID,
	[]coretask.Dependency{
		k8saudit.ManifestGeneratorTaskID.Ref(),
	},
	func(ctx context.Context, taskMode inspectioncore.InspectionTaskModeType) (k8scommon.NEGNameToResourceIdentityMap, error) {
		if taskMode == inspectioncore.TaskModeDryRun {
			return nil, nil
		}
		result := k8scommon.NEGNameToResourceIdentityMap{}
		resourceLogs := coretask.GetTaskResult(ctx, k8saudit.ManifestGeneratorTaskID.Ref())
		for _, group := range resourceLogs {
			if group.Resource.Type() != k8saudit.Resource {
				continue
			}
			if group.Resource.APIVersion != "networking.gke.io/v1beta1" || group.Resource.Kind != "servicenetworkendpointgroup" {
				continue
			}
			result[group.Resource.Name] = *group.Resource
		}
		return result, nil
	},
	coretask.ProvidesTag(k8scommon.TagNEGNamesDiscovery),
)

NEGToBackendServiceInventoryTask is the inventory task that provides aggregated NEG to BackendService mappings.

Functions

func Register

func Register(registry coreinspection.InspectionTaskRegistry) error

Register registers all googlecloudk8scommon inspection tasks to the registry.

Types

This section is empty.

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL