googlecloudclustercomposer_impl

package
v0.58.2 Latest Latest
Warning

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

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

Documentation

Index

Constants

This section is empty.

Variables

View Source
var AirflowDagProcessorManagerLogFilterTask = componentFilterTask(googlecloudclustercomposer_contract.AirflowDagProcessorManagerLogFilterTaskID, googlecloudclustercomposer_contract.ComposerLogsQueryTaskID.Ref(), "dag-processor-manager")
View Source
var AirflowDagProcessorManagerLogGrouperTask = inspectiontaskbase.NewLogGrouperTask(
	googlecloudclustercomposer_contract.AirflowDagProcessorManagerLogGrouperTaskID,
	googlecloudclustercomposer_contract.AirflowDagProcessorManagerLogSorterTaskID.Ref(),
	func(ctx context.Context, l *log.Log) string {
		fs, err := googlecloudclustercomposer_contract.ExtractComposer(l.NodeReader)
		if err != nil {
			return ""
		}
		if fs.SchedulerID != "" {
			return fs.SchedulerID
		}
		return fs.DagProcessorManagerID
	},
)

AirflowDagProcessorManagerLogGrouperTask groups Airflow DAG processor manager logs.

View Source
var AirflowDagProcessorManagerLogIngesterTask = inspectiontaskbase.NewGroupedLogIngesterTask(
	googlecloudclustercomposer_contract.AirflowDagProcessorManagerLogIngesterTaskID,
	&dagProcessorManagerLogIngester{},
)

AirflowDagProcessorManagerLogIngesterTask is the task that ingests Airflow DAG processor manager logs.

AirflowDagProcessorManagerLogSorterTask sorts Airflow DAG processor manager logs.

View Source
var AirflowDagProcessorManagerLogToTimelineMapperTask = inspectiontaskbase.NewLogToTimelineMapperTask(
	googlecloudclustercomposer_contract.AirflowDagProcessorManagerLogToTimelineMapperTaskID,
	&dagProcessorManagerTimelineMapper{
		targetLogType: googlecloudclustercomposer_contract.LogTypeManagedAirflowEnvironment,
		dagFilePath:   "/home/airflow/gcs/dags",
	},
)

AirflowDagProcessorManagerLogToTimelineMapperTask is the task that maps Airflow DAG processor manager logs to timeline events.

View Source
var AirflowOtherLogFilterTask = inspectiontaskbase.NewLogFilterTask(
	googlecloudclustercomposer_contract.AirflowOtherLogFilterTaskID,
	googlecloudclustercomposer_contract.ComposerLogsQueryTaskID.Ref(),
	func(ctx context.Context, l *log.Log) bool {
		component, err := googlecloudclustercomposer_contract.ExtractComposerComponent(l.NodeReader)
		if err != nil {
			return false
		}

		return component != "airflow-worker" && component != "airflow-scheduler" && component != "dag-processor-manager"
	},
)

AirflowOtherLogGrouperTask groups other Airflow logs.

AirflowOtherLogIngesterTask is the task that ingests other Airflow logs.

AirflowOtherLogToTimelineMapperTask is the task that maps other Airflow logs to timeline events.

View Source
var AirflowSchedulerLogFilterTask = componentFilterTask(googlecloudclustercomposer_contract.AirflowSchedulerLogFilterTaskID, googlecloudclustercomposer_contract.ComposerLogsQueryTaskID.Ref(), "airflow-scheduler")

AirflowSchedulerLogGrouperTask groups Airflow scheduler logs.

AirflowSchedulerLogIngesterTask is the task that ingests Airflow scheduler logs.

View Source
var AirflowSchedulerLogToTimelineMapperTask = inspectiontaskbase.NewLogToTimelineMapperTask(
	googlecloudclustercomposer_contract.AirflowSchedulerLogToTimelineMapperTaskID,
	&schedulerLogToTimelineMapper{},
)

AirflowSchedulerLogToTimelineMapperTask is the task that maps Airflow scheduler logs to timeline events.

AirflowWorkerLogGrouperTask groups Airflow worker logs.

AirflowWorkerLogIngesterTask is the task that ingests Airflow worker logs.

AirflowWorkerLogToTimelineMapperTask is the task that maps Airflow worker logs to timeline events.

View Source
var AutocompleteComposerClusterNamesTask = inspectiontaskbase.NewGlobalCachedTask(googlecloudclustercomposer_contract.AutocompleteComposerClusterNamesTaskID, []taskid.UntypedTaskReference{
	googlecloudclustercomposer_contract.ComposerEnvironmentClusterFinderTaskID.Ref(),
	googlecloudcommon_contract.InputProjectIdTaskID.Ref(),
	googlecloudcommon_contract.InputLocationsTaskID.Ref(),
	googlecloudclustercomposer_contract.InputComposerEnvironmentNameTaskID.Ref(),
	googlecloudclustercomposer_contract.AutocompleteComposerEnvironmentIdentityTaskID.Ref(),
	googlecloudcommon_contract.InputStartTimeTaskID.Ref(),
	googlecloudcommon_contract.InputEndTimeTaskID.Ref(),
}, func(ctx context.Context, prevValue inspectiontaskbase.CacheableTaskResult[*inspectioncore_contract.AutocompleteResult[googlecloudk8scommon_contract.GoogleCloudClusterIdentity]]) (inspectiontaskbase.CacheableTaskResult[*inspectioncore_contract.AutocompleteResult[googlecloudk8scommon_contract.GoogleCloudClusterIdentity]], error) {

	projectID := coretask.GetTaskResult(ctx, googlecloudcommon_contract.InputProjectIdTaskID.Ref())
	environment := coretask.GetTaskResult(ctx, googlecloudclustercomposer_contract.InputComposerEnvironmentNameTaskID.Ref())
	location := coretask.GetTaskResult(ctx, googlecloudcommon_contract.InputLocationsTaskID.Ref())
	startTime := coretask.GetTaskResult(ctx, googlecloudcommon_contract.InputStartTimeTaskID.Ref())
	endTime := coretask.GetTaskResult(ctx, googlecloudcommon_contract.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_contract.AutocompleteResult[googlecloudk8scommon_contract.GoogleCloudClusterIdentity]]{
			DependencyDigest: dependencyDigest,
			Value: &inspectioncore_contract.AutocompleteResult[googlecloudk8scommon_contract.GoogleCloudClusterIdentity]{
				Values: []googlecloudk8scommon_contract.GoogleCloudClusterIdentity{},
				Error:  "Project ID or Composer environment name is empty",
			},
		}, nil
	}

	if location == "" {
		return inspectiontaskbase.CacheableTaskResult[*inspectioncore_contract.AutocompleteResult[googlecloudk8scommon_contract.GoogleCloudClusterIdentity]]{
			DependencyDigest: dependencyDigest,
			Value: &inspectioncore_contract.AutocompleteResult[googlecloudk8scommon_contract.GoogleCloudClusterIdentity]{
				Values: []googlecloudk8scommon_contract.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, googlecloudclustercomposer_contract.ComposerEnvironmentClusterFinderTaskID.Ref())
	clusterNames, err := clusterFinder.GetGKEClusterNames(ctx, projectID, location, environment, startTime, endTime)
	if err != nil {
		if errors.Is(err, googlecloudclustercomposer_contract.ErrEnvironmentClusterNotFound) {
			return inspectiontaskbase.CacheableTaskResult[*inspectioncore_contract.AutocompleteResult[googlecloudk8scommon_contract.GoogleCloudClusterIdentity]]{
				DependencyDigest: dependencyDigest,
				Value: &inspectioncore_contract.AutocompleteResult[googlecloudk8scommon_contract.GoogleCloudClusterIdentity]{
					Values: []googlecloudk8scommon_contract.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_contract.AutocompleteResult[googlecloudk8scommon_contract.GoogleCloudClusterIdentity]]{
			DependencyDigest: dependencyDigest,
			Value: &inspectioncore_contract.AutocompleteResult[googlecloudk8scommon_contract.GoogleCloudClusterIdentity]{
				Values: []googlecloudk8scommon_contract.GoogleCloudClusterIdentity{},
				Error:  "Failed to fetch the list GKE cluster. Please confirm if the Project ID is correct, or retry later",
			},
		}, nil
	}

	identities := make([]googlecloudk8scommon_contract.GoogleCloudClusterIdentity, len(clusterNames))
	for i, clusterName := range clusterNames {
		identities[i] = googlecloudk8scommon_contract.GoogleCloudClusterIdentity{
			ClusterName: clusterName,
			ProjectID:   projectID,
			Location:    location,
		}
	}

	return inspectiontaskbase.CacheableTaskResult[*inspectioncore_contract.AutocompleteResult[googlecloudk8scommon_contract.GoogleCloudClusterIdentity]]{
		DependencyDigest: dependencyDigest,
		Value: &inspectioncore_contract.AutocompleteResult[googlecloudk8scommon_contract.GoogleCloudClusterIdentity]{
			Values: identities,
		},
	}, nil
},
	coretask.WithSelectionPriority(1000),
)

AutocompleteComposerClusterNamesTask is an implementation for googlecloudk8scommon_contract.AutocompleteClusterNamesTaskID the task returns GKE cluster name where the provided Composer environment is running.

View Source
var AutocompleteComposerComponentsTask = inspectiontaskbase.NewGlobalCachedTask(googlecloudclustercomposer_contract.AutocompleteComposerComponentsTaskID, []taskid.UntypedTaskReference{
	googlecloudclustercomposer_contract.ClusterIdentityTaskID.GetUntypedReference(),
	googlecloudcommon_contract.InputStartTimeTaskID.Ref(),
	googlecloudcommon_contract.InputEndTimeTaskID.Ref(),
	googlecloudclustercomposer_contract.InputComposerEnvironmentNameTaskID.Ref(),
	googlecloudcommon_contract.APIClientFactoryTaskID.Ref(),
	googlecloudcommon_contract.APIClientCallOptionsInjectorTaskID.Ref(),
}, func(ctx context.Context, prevValue inspectiontaskbase.CacheableTaskResult[*inspectioncore_contract.AutocompleteResult[string]]) (inspectiontaskbase.CacheableTaskResult[*inspectioncore_contract.AutocompleteResult[string]], error) {
	clusterIdentity := coretask.GetTaskResult(ctx, googlecloudclustercomposer_contract.ClusterIdentityTaskID.Ref())
	projectID := clusterIdentity.ProjectID
	location := clusterIdentity.Location

	startTime := coretask.GetTaskResult(ctx, googlecloudcommon_contract.InputStartTimeTaskID.Ref())
	endTime := coretask.GetTaskResult(ctx, googlecloudcommon_contract.InputEndTimeTaskID.Ref())
	environmentName := coretask.GetTaskResult(ctx, googlecloudclustercomposer_contract.InputComposerEnvironmentNameTaskID.Ref())
	cf := coretask.GetTaskResult(ctx, googlecloudcommon_contract.APIClientFactoryTaskID.Ref())
	optionInjector := coretask.GetTaskResult(ctx, googlecloudcommon_contract.APIClientCallOptionsInjectorTaskID.Ref())

	currentDigest := fmt.Sprintf("%s-%s-%s-%s-%d-%d", projectID, location, environmentName, "logging.googleapis.com/log_entry_count", startTime.Unix(), endTime.Unix())
	if currentDigest == prevValue.DependencyDigest {
		return prevValue, nil
	}

	if projectID == "" || environmentName == "" || location == "" {
		return inspectiontaskbase.CacheableTaskResult[*inspectioncore_contract.AutocompleteResult[string]]{
			Value: &inspectioncore_contract.AutocompleteResult[string]{
				Values: []string{},
				Hint:   "Components are suggested after the project ID, location, and environment name are provided.",
			},
			DependencyDigest: currentDigest,
		}, nil
	}

	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(`resource.type = "cloud_composer_environment" AND metric.type = "logging.googleapis.com/log_entry_count" AND resource.labels.environment_name = "%s" AND resource.labels.location = "%s"`, environmentName, location)

	errorString := ""
	hintString := ""
	metricsLabels, err := googlecloud.QueryResourceLabelsFromMetrics(ctx, client, projectID, filter, startTime, endTime, []string{"metric.label.log"})
	if err != nil {
		errorString = err.Error()
	}

	componentsMap := make(map[string]struct{})
	for _, labels := range metricsLabels {
		if logName, ok := labels["log"]; ok && logName != "" {
			componentsMap[logName] = struct{}{}
		}
	}

	components := make([]string, 0, len(componentsMap))
	for comp := range componentsMap {
		components = append(components, comp)
	}
	sort.Strings(components)

	if hintString == "" && errorString == "" && len(components) == 0 {
		hintString = "No components found for the specified environment and time range."
	}

	return inspectiontaskbase.CacheableTaskResult[*inspectioncore_contract.AutocompleteResult[string]]{
		DependencyDigest: currentDigest,
		Value: &inspectioncore_contract.AutocompleteResult[string]{
			Values: components,
			Error:  errorString,
			Hint:   hintString,
		},
	}, nil
})
View Source
var AutocompleteComposerEnvironmentIdentityTask = inspectiontaskbase.NewGlobalCachedTask(googlecloudclustercomposer_contract.AutocompleteComposerEnvironmentIdentityTaskID, []taskid.UntypedTaskReference{
	googlecloudcommon_contract.InputProjectIdTaskID.Ref(),
	googlecloudcommon_contract.InputStartTimeTaskID.Ref(),
	googlecloudcommon_contract.InputEndTimeTaskID.Ref(),
	googlecloudcommon_contract.APIClientFactoryTaskID.Ref(),
	googlecloudcommon_contract.APIClientCallOptionsInjectorTaskID.Ref(),
}, func(ctx context.Context, prevValue inspectiontaskbase.CacheableTaskResult[*inspectioncore_contract.AutocompleteResult[googlecloudclustercomposer_contract.ComposerEnvironmentIdentity]]) (inspectiontaskbase.CacheableTaskResult[*inspectioncore_contract.AutocompleteResult[googlecloudclustercomposer_contract.ComposerEnvironmentIdentity]], error) {
	projectID := coretask.GetTaskResult(ctx, googlecloudcommon_contract.InputProjectIdTaskID.Ref())
	startTime := coretask.GetTaskResult(ctx, googlecloudcommon_contract.InputStartTimeTaskID.Ref())
	endTime := coretask.GetTaskResult(ctx, googlecloudcommon_contract.InputEndTimeTaskID.Ref())
	cf := coretask.GetTaskResult(ctx, googlecloudcommon_contract.APIClientFactoryTaskID.Ref())
	optionInjector := coretask.GetTaskResult(ctx, googlecloudcommon_contract.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_contract.AutocompleteResult[googlecloudclustercomposer_contract.ComposerEnvironmentIdentity]]{
			Value: &inspectioncore_contract.AutocompleteResult[googlecloudclustercomposer_contract.ComposerEnvironmentIdentity]{
				Values: []googlecloudclustercomposer_contract.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([]googlecloudclustercomposer_contract.ComposerEnvironmentIdentity, 0, len(metricsLabels))
	for _, labels := range metricsLabels {
		envName := labels["environment_name"]
		location := labels["location"]
		if envName != "" && location != "" {
			identities = append(identities, googlecloudclustercomposer_contract.ComposerEnvironmentIdentity{
				ProjectID:       projectID,
				Location:        location,
				EnvironmentName: envName,
			})
		}
	}

	return inspectiontaskbase.CacheableTaskResult[*inspectioncore_contract.AutocompleteResult[googlecloudclustercomposer_contract.ComposerEnvironmentIdentity]]{
		DependencyDigest: currentDigest,
		Value: &inspectioncore_contract.AutocompleteResult[googlecloudclustercomposer_contract.ComposerEnvironmentIdentity]{
			Values: identities,
			Error:  errorString,
			Hint:   hintString,
		},
	}, nil
})

AutocompleteComposerEnvironmentIdentityTask is the task that autocompletes composer environment identities.

View Source
var AutocompleteLocationForComposerEnvironmentTask = inspectiontaskbase.NewGlobalCachedTask(googlecloudclustercomposer_contract.AutocompleteLocationForComposerEnvironmentTaskID, []taskid.UntypedTaskReference{
	googlecloudclustercomposer_contract.AutocompleteComposerEnvironmentIdentityTaskID.Ref(),
	googlecloudcommon_contract.InputProjectIdTaskID.Ref(),
	googlecloudclustercomposer_contract.InputComposerEnvironmentNameTaskID.Ref(),
	googlecloudcommon_contract.InputStartTimeTaskID.Ref(),
	googlecloudcommon_contract.InputEndTimeTaskID.Ref(),
}, func(ctx context.Context, prevValue inspectiontaskbase.CacheableTaskResult[*inspectioncore_contract.AutocompleteResult[string]]) (inspectiontaskbase.CacheableTaskResult[*inspectioncore_contract.AutocompleteResult[string]], error) {
	projectID := coretask.GetTaskResult(ctx, googlecloudcommon_contract.InputProjectIdTaskID.Ref())
	environmentName := coretask.GetTaskResult(ctx, googlecloudclustercomposer_contract.InputComposerEnvironmentNameTaskID.Ref())
	startTime := coretask.GetTaskResult(ctx, googlecloudcommon_contract.InputStartTimeTaskID.Ref())
	endTime := coretask.GetTaskResult(ctx, googlecloudcommon_contract.InputEndTimeTaskID.Ref())
	identities := coretask.GetTaskResult(ctx, googlecloudclustercomposer_contract.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_contract.AutocompleteResult[string]]{
			Value: &inspectioncore_contract.AutocompleteResult[string]{
				Values: []string{},
				Error:  "",
				Hint:   "Locations are suggested after the project ID is provided.",
			},
			DependencyDigest: currentDigest,
		}, nil
	}

	if environmentName == "" {
		return inspectiontaskbase.CacheableTaskResult[*inspectioncore_contract.AutocompleteResult[string]]{
			Value: &inspectioncore_contract.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_contract.AutocompleteResult[string]]{
			Value: &inspectioncore_contract.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_contract.AutocompleteResult[string]]{
		Value: &inspectioncore_contract.AutocompleteResult[string]{
			Values: locations,
			Error:  "",
			Hint:   identities.Hint,
		},
		DependencyDigest: currentDigest,
	}, nil
},
	coretask.WithSelectionPriority(1000),
)

ComposerEnvironmentClusterFinderTask injects ComposerEnvironmentClusterFinder implementation.

ComposerEnvironmentListFetcherTask injects ComposerEnvironmentListFetcher implementation.

View Source
var ComposerLogsQueryTask = googlecloudcommon_contract.NewStructuredListLogEntriesTask(&composerListLogEntriesTaskSetting{
	taskId:    googlecloudclustercomposer_contract.ComposerLogsQueryTaskID,
	queryName: "Composer Environment Logs",
})

ComposerLogsQueryTask defines a task that gathers logs from Cloud Logging for multiple Composer components.

View Source
var ComposerLogsTailTask = coretask.NewTask(
	googlecloudclustercomposer_contract.ComposerLogsTailTaskID,
	[]taskid.UntypedTaskReference{
		googlecloudclustercomposer_contract.AirflowWorkerLogToTimelineMapperTaskID.Ref(),
		googlecloudclustercomposer_contract.AirflowSchedulerLogToTimelineMapperTaskID.Ref(),
		googlecloudclustercomposer_contract.AirflowDagProcessorManagerLogToTimelineMapperTaskID.Ref(),
		googlecloudclustercomposer_contract.AirflowOtherLogToTimelineMapperTaskID.Ref(),
	},
	func(ctx context.Context) (struct{}, error) {
		return struct{}{}, nil
	},
	inspectioncore_contract.FeatureTaskLabel(
		"Managed Service for Apache Airflow Logs",
		"Gather Managed Service for Apache Airflow logs, including airflow-worker, airflow-scheduler, and airflow-dag-processor-manager, to visualize general environment operations on timelines.",
		2500,
		true,
	),
)
View Source
var InputComposerComponentsTask = formtask.NewSetFormTaskBuilder(googlecloudclustercomposer_contract.InputComposerComponentsTaskID, googlecloudcommon_contract.FormBasePriority+3000, "Composer Components").
	WithDependencies([]taskid.UntypedTaskReference{googlecloudclustercomposer_contract.AutocompleteComposerComponentsTaskID.Ref()}).
	WithDefaultValueConstant([]string{"@any"}, true).
	WithAllowAddAll(false).
	WithAllowRemoveAll(false).
	WithAllowCustomValue(false).
	WithDescription(`Select which Composer V3 components to fetch logs from.`).
	WithOptionsFunc(func(ctx context.Context, previousValues []string) ([]inspectionmetadata.SetParameterFormFieldOptionItem, error) {
		autocompleteResult := coretask.GetTaskResult(ctx, googlecloudclustercomposer_contract.AutocompleteComposerComponentsTaskID.Ref())

		var options []inspectionmetadata.SetParameterFormFieldOptionItem
		options = append(options, inspectionmetadata.SetParameterFormFieldOptionItem{
			ID: "@any",
		})
		if autocompleteResult != nil {
			for _, comp := range autocompleteResult.Values {
				options = append(options, inspectionmetadata.SetParameterFormFieldOptionItem{
					ID: comp,
				})
			}
		}
		return options, nil
	}).
	WithHintFunc(func(ctx context.Context, value []string, convertedValue any) (string, inspectionmetadata.ParameterHintType, error) {
		autocompleteResult := coretask.GetTaskResult(ctx, googlecloudclustercomposer_contract.AutocompleteComposerComponentsTaskID.Ref())
		if autocompleteResult != nil {
			if autocompleteResult.Error != "" {
				return autocompleteResult.Error, inspectionmetadata.Error, nil
			}
			if autocompleteResult.Hint != "" {
				return autocompleteResult.Hint, inspectionmetadata.Info, nil
			}
		}
		return "", inspectionmetadata.None, nil
	}).
	WithConverter(func(ctx context.Context, value []string) ([]string, error) {
		return value, nil
	}).
	Build()
View Source
var InputComposerEnvironmentNameTask = formtask.NewTextFormTaskBuilder(googlecloudclustercomposer_contract.InputComposerEnvironmentNameTaskID, googlecloudcommon_contract.PriorityForResourceIdentifierGroup+4400, "Composer Environment Name").WithDependencies(
	[]taskid.UntypedTaskReference{googlecloudclustercomposer_contract.AutocompleteComposerEnvironmentIdentityTaskID.Ref()},
).WithSuggestionsFunc(func(ctx context.Context, value string, previousValues []string) ([]string, error) {
	environments := coretask.GetTaskResult(ctx, googlecloudclustercomposer_contract.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 GenerateComposerLogsQuery added in v0.58.2

func GenerateComposerLogsQuery(projectID, location, environmentName string, selectedComponents []string) string

GenerateComposerLogsQuery generates a query for Composer environment logs.

func GenerateComposerLogsStructuredQuery added in v0.58.2

func GenerateComposerLogsStructuredQuery(projectID, location, environmentName string, selectedComponents []string) *logestimator.StructuredLogQuery

GenerateComposerLogsStructuredQuery generates a structured query for Composer environment logs.

func Register

func Register(registry coreinspection.InspectionTaskRegistry) error

Register registers all googlecloudclustercomposer inspection tasks to the registry.

Types

type DagProcessorState added in v0.52.8

type DagProcessorState struct {
	Reader *logutil.TabulateReader
}

DagProcessorState retains the parsing state using TabulateReader.

Jump to

Keyboard shortcuts

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