Documentation
¶
Index ¶
Constants ¶
This section is empty.
Variables ¶
View Source
var AutocompleteComposerClusterNamesTask = inspectiontaskbase.NewGlobalCachedTask(composercluster.AutocompleteComposerClusterNamesTaskID, []coretask.Dependency{ composercluster.ComposerEnvironmentClusterFinderTaskID.Ref(), gcpcommon.InputProjectIdTaskID.Ref(), gcpcommon.InputLocationsTaskID.Ref(), composercluster.InputComposerEnvironmentNameTaskID.Ref(), composercluster.AutocompleteComposerEnvironmentIdentityTaskID.Ref(), gcpcommon.InputStartTimeTaskID.Ref(), gcpcommon.InputEndTimeTaskID.Ref(), }, func(ctx context.Context, prevValue inspectiontaskbase.CacheableTaskResult[*inspectioncore.AutocompleteResult[k8scommon.GoogleCloudClusterIdentity]]) (inspectiontaskbase.CacheableTaskResult[*inspectioncore.AutocompleteResult[k8scommon.GoogleCloudClusterIdentity]], error) { projectID := coretask.GetTaskResult(ctx, gcpcommon.InputProjectIdTaskID.Ref()) environment := coretask.GetTaskResult(ctx, composercluster.InputComposerEnvironmentNameTaskID.Ref()) location := coretask.GetTaskResult(ctx, gcpcommon.InputLocationsTaskID.Ref()) startTime := coretask.GetTaskResult(ctx, gcpcommon.InputStartTimeTaskID.Ref()) endTime := coretask.GetTaskResult(ctx, gcpcommon.InputEndTimeTaskID.Ref()) dependencyDigest := fmt.Sprintf("%s-%s-%s-%d-%d", projectID, environment, location, startTime.Unix(), endTime.Unix()) isWIP := projectID == "" || environment == "" if isWIP { return inspectiontaskbase.CacheableTaskResult[*inspectioncore.AutocompleteResult[k8scommon.GoogleCloudClusterIdentity]]{ DependencyDigest: dependencyDigest, Value: &inspectioncore.AutocompleteResult[k8scommon.GoogleCloudClusterIdentity]{ Values: []k8scommon.GoogleCloudClusterIdentity{}, Error: "Project ID or Composer environment name is empty", }, }, nil } if location == "" { return inspectiontaskbase.CacheableTaskResult[*inspectioncore.AutocompleteResult[k8scommon.GoogleCloudClusterIdentity]]{ DependencyDigest: dependencyDigest, Value: &inspectioncore.AutocompleteResult[k8scommon.GoogleCloudClusterIdentity]{ Values: []k8scommon.GoogleCloudClusterIdentity{}, Error: "", Hint: "Cluster names are suggested after the location is provided.", }, }, nil } if environment != "" && dependencyDigest == prevValue.DependencyDigest { return prevValue, nil } clusterFinder := coretask.GetTaskResult(ctx, composercluster.ComposerEnvironmentClusterFinderTaskID.Ref()) clusterNames, err := clusterFinder.GetGKEClusterNames(ctx, projectID, location, environment, startTime, endTime) if err != nil { if errors.Is(err, composercluster.ErrEnvironmentClusterNotFound) { return inspectiontaskbase.CacheableTaskResult[*inspectioncore.AutocompleteResult[k8scommon.GoogleCloudClusterIdentity]]{ DependencyDigest: dependencyDigest, Value: &inspectioncore.AutocompleteResult[k8scommon.GoogleCloudClusterIdentity]{ Values: []k8scommon.GoogleCloudClusterIdentity{}, Error: `Not found. 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. Note: Composer 3 is not running on your GKE cluster. Please remove all Kubernetes/GKE queries from the previous section.`, }, }, nil } return inspectiontaskbase.CacheableTaskResult[*inspectioncore.AutocompleteResult[k8scommon.GoogleCloudClusterIdentity]]{ DependencyDigest: dependencyDigest, Value: &inspectioncore.AutocompleteResult[k8scommon.GoogleCloudClusterIdentity]{ Values: []k8scommon.GoogleCloudClusterIdentity{}, Error: "Failed to fetch the list GKE cluster. Please confirm if the Project ID is correct, or retry later", }, }, nil } identities := make([]k8scommon.GoogleCloudClusterIdentity, len(clusterNames)) for i, clusterName := range clusterNames { identities[i] = k8scommon.GoogleCloudClusterIdentity{ ClusterName: clusterName, ProjectID: projectID, Location: location, } } return inspectiontaskbase.CacheableTaskResult[*inspectioncore.AutocompleteResult[k8scommon.GoogleCloudClusterIdentity]]{ DependencyDigest: dependencyDigest, Value: &inspectioncore.AutocompleteResult[k8scommon.GoogleCloudClusterIdentity]{ Values: identities, }, }, nil }, coretask.WithSelectionPriority(1000), )
AutocompleteComposerClusterNamesTask is an implementation for k8scommon.AutocompleteClusterNamesTaskID the task returns GKE cluster name where the provided Composer environment is running.
View Source
var AutocompleteComposerEnvironmentIdentityTask = inspectiontaskbase.NewGlobalCachedTask(composercluster.AutocompleteComposerEnvironmentIdentityTaskID, []coretask.Dependency{ gcpcommon.InputProjectIdTaskID.Ref(), gcpcommon.InputStartTimeTaskID.Ref(), gcpcommon.InputEndTimeTaskID.Ref(), gcpcommon.APIClientFactoryTaskID.Ref(), gcpcommon.APIClientCallOptionsInjectorTaskID.Ref(), }, func(ctx context.Context, prevValue inspectiontaskbase.CacheableTaskResult[*inspectioncore.AutocompleteResult[composercluster.ComposerEnvironmentIdentity]]) (inspectiontaskbase.CacheableTaskResult[*inspectioncore.AutocompleteResult[composercluster.ComposerEnvironmentIdentity]], error) { projectID := coretask.GetTaskResult(ctx, gcpcommon.InputProjectIdTaskID.Ref()) startTime := coretask.GetTaskResult(ctx, gcpcommon.InputStartTimeTaskID.Ref()) endTime := coretask.GetTaskResult(ctx, gcpcommon.InputEndTimeTaskID.Ref()) cf := coretask.GetTaskResult(ctx, gcpcommon.APIClientFactoryTaskID.Ref()) optionInjector := coretask.GetTaskResult(ctx, gcpcommon.APIClientCallOptionsInjectorTaskID.Ref()) currentDigest := fmt.Sprintf("%s-%d-%d", projectID, startTime.Unix(), endTime.Unix()) if currentDigest == prevValue.DependencyDigest { return prevValue, nil } if projectID == "" { return inspectiontaskbase.CacheableTaskResult[*inspectioncore.AutocompleteResult[composercluster.ComposerEnvironmentIdentity]]{ Value: &inspectioncore.AutocompleteResult[composercluster.ComposerEnvironmentIdentity]{ Values: []composercluster.ComposerEnvironmentIdentity{}, Error: "", Hint: "Composer environments 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 environment 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 := `metric.type="composer.googleapis.com/environment/healthy" AND resource.type="cloud_composer_environment"` metricsLabels, err := googlecloud.QueryResourceLabelsFromMetrics(ctx, client, projectID, filter, startTime, endTime, []string{"resource.label.environment_name", "resource.label.location"}) if err != nil { errorString = err.Error() } if hintString == "" && errorString == "" && len(metricsLabels) == 0 { hintString = fmt.Sprintf("No Composer environments 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 environment name.", startTime.Format(time.RFC3339), endTime.Format(time.RFC3339)) } identities := make([]composercluster.ComposerEnvironmentIdentity, 0, len(metricsLabels)) for _, labels := range metricsLabels { envName := labels["environment_name"] location := labels["location"] if envName != "" && location != "" { identities = append(identities, composercluster.ComposerEnvironmentIdentity{ ProjectID: projectID, Location: location, EnvironmentName: envName, }) } } return inspectiontaskbase.CacheableTaskResult[*inspectioncore.AutocompleteResult[composercluster.ComposerEnvironmentIdentity]]{ DependencyDigest: currentDigest, Value: &inspectioncore.AutocompleteResult[composercluster.ComposerEnvironmentIdentity]{ Values: identities, Error: errorString, Hint: hintString, }, }, nil })
AutocompleteComposerEnvironmentIdentityTask is the task that autocompletes composer environment identities.
View Source
var AutocompleteLocationForComposerEnvironmentTask = inspectiontaskbase.NewGlobalCachedTask(composercluster.AutocompleteLocationForComposerEnvironmentTaskID, []coretask.Dependency{ composercluster.AutocompleteComposerEnvironmentIdentityTaskID.Ref(), gcpcommon.InputProjectIdTaskID.Ref(), composercluster.InputComposerEnvironmentNameTaskID.Ref(), gcpcommon.InputStartTimeTaskID.Ref(), gcpcommon.InputEndTimeTaskID.Ref(), }, func(ctx context.Context, prevValue inspectiontaskbase.CacheableTaskResult[*inspectioncore.AutocompleteResult[string]]) (inspectiontaskbase.CacheableTaskResult[*inspectioncore.AutocompleteResult[string]], error) { projectID := coretask.GetTaskResult(ctx, gcpcommon.InputProjectIdTaskID.Ref()) environmentName := coretask.GetTaskResult(ctx, composercluster.InputComposerEnvironmentNameTaskID.Ref()) startTime := coretask.GetTaskResult(ctx, gcpcommon.InputStartTimeTaskID.Ref()) endTime := coretask.GetTaskResult(ctx, gcpcommon.InputEndTimeTaskID.Ref()) identities := coretask.GetTaskResult(ctx, composercluster.AutocompleteComposerEnvironmentIdentityTaskID.Ref()) currentDigest := fmt.Sprintf("%s-%s-%d-%d", projectID, environmentName, 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 are suggested after the project ID is provided.", }, DependencyDigest: currentDigest, }, nil } if environmentName == "" { return inspectiontaskbase.CacheableTaskResult[*inspectioncore.AutocompleteResult[string]]{ Value: &inspectioncore.AutocompleteResult[string]{ Values: []string{}, Error: "", Hint: "Locations are suggested after the environment name is provided.", }, DependencyDigest: currentDigest, }, nil } if identities.Error != "" { return inspectiontaskbase.CacheableTaskResult[*inspectioncore.AutocompleteResult[string]]{ Value: &inspectioncore.AutocompleteResult[string]{ Values: []string{}, Error: identities.Error, Hint: identities.Hint, }, DependencyDigest: currentDigest, }, nil } locationsMap := make(map[string]struct{}) for _, identity := range identities.Values { if identity.EnvironmentName == environmentName { locationsMap[identity.Location] = struct{}{} } } locations := make([]string, 0, len(locationsMap)) for location := range locationsMap { locations = append(locations, location) } return inspectiontaskbase.CacheableTaskResult[*inspectioncore.AutocompleteResult[string]]{ Value: &inspectioncore.AutocompleteResult[string]{ Values: locations, Error: "", Hint: identities.Hint, }, DependencyDigest: currentDigest, }, nil }, coretask.WithSelectionPriority(1000), )
View Source
var ClusterIdentityAliasTask = coretask.NewAliasTask( composercluster.ClusterIdentityTaskID, k8scommon.ClusterIdentityTaskID.Ref(), )
View Source
var ComposerEnvironmentClusterFinderTask = coretask.NewTask( composercluster.ComposerEnvironmentClusterFinderTaskID, []coretask.Dependency{ gcpcommon.APIClientFactoryTaskID.Ref(), gcpcommon.APIClientCallOptionsInjectorTaskID.Ref(), }, func(ctx context.Context) (composercluster.ComposerEnvironmentClusterFinder, error) { cf := coretask.GetTaskResult(ctx, gcpcommon.APIClientFactoryTaskID.Ref()) injector := coretask.GetTaskResult(ctx, gcpcommon.APIClientCallOptionsInjectorTaskID.Ref()) return composercluster.NewEnvironmentClusterFinder(cf, injector), nil }, )
ComposerEnvironmentClusterFinderTask injects ComposerEnvironmentClusterFinder implementation.
View Source
var InputComposerEnvironmentNameTask = formtask.NewTextFormTaskBuilder(composercluster.InputComposerEnvironmentNameTaskID, gcpcommon.PriorityForResourceIdentifierGroup+4400, "Composer Environment Name").WithDependencies( []coretask.Dependency{composercluster.AutocompleteComposerEnvironmentIdentityTaskID.Ref()}, ).WithDefaultValueFunc(func(ctx context.Context, previousValues []string) (string, error) { environments := coretask.GetTaskResult(ctx, composercluster.AutocompleteComposerEnvironmentIdentityTaskID.Ref()) if len(previousValues) > 0 && slices.ContainsFunc(environments.Values, func(env composercluster.ComposerEnvironmentIdentity) bool { return env.EnvironmentName == previousValues[0] }) { return previousValues[0], nil } if len(environments.Values) == 0 { return "", nil } return environments.Values[0].EnvironmentName, nil }).WithSuggestionsFunc(func(ctx context.Context, value string, previousValues []string) ([]string, error) { environments := coretask.GetTaskResult(ctx, composercluster.AutocompleteComposerEnvironmentIdentityTaskID.Ref()) if environments.Error != "" { return []string{}, nil } environmentNames := make([]string, len(environments.Values)) for i, env := range environments.Values { environmentNames[i] = env.EnvironmentName } return common.SortForAutocomplete(value, environmentNames), nil }).Build()
InputComposerEnvironmentNameTask is the task that inputs composer environment name.
Functions ¶
func Register ¶
func Register(registry coreinspection.InspectionTaskRegistry) error
Register registers all composer cluster inspection tasks to the registry.
Types ¶
This section is empty.
Click to show internal directories.
Click to hide internal directories.