Documentation
¶
Index ¶
Constants ¶
This section is empty.
Variables ¶
var AirflowDagProcessorManagerLogFilterTask = componentFilterTask(googlecloudclustercomposer_contract.AirflowDagProcessorManagerLogFilterTaskID, googlecloudclustercomposer_contract.ComposerLogsFieldSetReadTaskID.Ref(), "dag-processor-manager")
var AirflowDagProcessorManagerLogGrouperTask = inspectiontaskbase.NewLogGrouperTask( googlecloudclustercomposer_contract.AirflowDagProcessorManagerLogGrouperTaskID, googlecloudclustercomposer_contract.AirflowDagProcessorManagerLogSorterTaskID.Ref(), func(ctx context.Context, l *log.Log) string { fs, err := log.GetFieldSet(l, &googlecloudclustercomposer_contract.ComposerFieldSet{}) if err != nil { return "" } if fs.SchedulerID != "" { return fs.SchedulerID } return fs.DagProcessorManagerID }, )
var AirflowDagProcessorManagerLogIngesterTask = inspectiontaskbase.NewLogIngesterTask( googlecloudclustercomposer_contract.AirflowDagProcessorManagerLogIngesterTaskID, googlecloudclustercomposer_contract.AirflowDagProcessorManagerLogFilterTaskID.Ref(), )
var AirflowDagProcessorManagerLogSorterTask = inspectiontaskbase.NewLogSorterByTimeTask( googlecloudclustercomposer_contract.AirflowDagProcessorManagerLogSorterTaskID, googlecloudclustercomposer_contract.AirflowDagProcessorManagerLogFilterTaskID.Ref(), )
var AirflowDagProcessorManagerLogToTimelineMapperTask = inspectiontaskbase.NewLogToTimelineMapperTask[*DagProcessorState]( googlecloudclustercomposer_contract.AirflowDagProcessorManagerLogToTimelineMapperTaskID, &airflowDagProcessorManagerLogToTimelineMapperSetting{ targetLogType: enum.LogTypeComposerEnvironment, dagFilePath: "/home/airflow/gcs/dags", }, )
var AirflowOtherLogFilterTask = inspectiontaskbase.NewLogFilterTask( googlecloudclustercomposer_contract.AirflowOtherLogFilterTaskID, googlecloudclustercomposer_contract.ComposerLogsFieldSetReadTaskID.Ref(), func(ctx context.Context, l *log.Log) bool { fs, err := log.GetFieldSet(l, &googlecloudclustercomposer_contract.ComposerFieldSet{}) if err != nil { return false } return fs.Component != "airflow-worker" && fs.Component != "airflow-scheduler" && fs.Component != "dag-processor-manager" }, )
var AirflowOtherLogGrouperTask = inspectiontaskbase.NewLogGrouperTask( googlecloudclustercomposer_contract.AirflowOtherLogGrouperTaskID, googlecloudclustercomposer_contract.AirflowOtherLogFilterTaskID.Ref(), func(ctx context.Context, l *log.Log) string { return "" }, )
var AirflowOtherLogIngesterTask = inspectiontaskbase.NewLogIngesterTask( googlecloudclustercomposer_contract.AirflowOtherLogIngesterTaskID, googlecloudclustercomposer_contract.AirflowOtherLogFilterTaskID.Ref(), )
var AirflowOtherLogToTimelineMapperTask = inspectiontaskbase.NewLogToTimelineMapperTask[struct{}]( googlecloudclustercomposer_contract.AirflowOtherLogToTimelineMapperTaskID, &airflowOtherLogToTimelineMapperSetting{}, )
var AirflowSchedulerLogFilterTask = componentFilterTask(googlecloudclustercomposer_contract.AirflowSchedulerLogFilterTaskID, googlecloudclustercomposer_contract.ComposerLogsFieldSetReadTaskID.Ref(), "airflow-scheduler")
var AirflowSchedulerLogGrouperTask = inspectiontaskbase.NewLogGrouperTask( googlecloudclustercomposer_contract.AirflowSchedulerLogGrouperTaskID, googlecloudclustercomposer_contract.AirflowSchedulerLogFilterTaskID.Ref(), func(ctx context.Context, l *log.Log) string { return "" }, )
var AirflowSchedulerLogIngesterTask = inspectiontaskbase.NewLogIngesterTask( googlecloudclustercomposer_contract.AirflowSchedulerLogIngesterTaskID, googlecloudclustercomposer_contract.AirflowSchedulerLogFilterTaskID.Ref(), )
var AirflowSchedulerLogToTimelineMapperTask = inspectiontaskbase.NewLogToTimelineMapperTask[struct{}]( googlecloudclustercomposer_contract.AirflowSchedulerLogToTimelineMapperTaskID, &airflowSchedulerLogToTimelineMapperSetting{ targetLogType: enum.LogTypeComposerEnvironment, }, )
var AirflowWorkerLogFilterTask = componentFilterTask(googlecloudclustercomposer_contract.AirflowWorkerLogFilterTaskID, googlecloudclustercomposer_contract.ComposerLogsFieldSetReadTaskID.Ref(), "airflow-worker")
var AirflowWorkerLogGrouperTask = inspectiontaskbase.NewLogGrouperTask( googlecloudclustercomposer_contract.AirflowWorkerLogGrouperTaskID, googlecloudclustercomposer_contract.AirflowWorkerLogFilterTaskID.Ref(), func(ctx context.Context, l *log.Log) string { return "" }, )
var AirflowWorkerLogIngesterTask = inspectiontaskbase.NewLogIngesterTask( googlecloudclustercomposer_contract.AirflowWorkerLogIngesterTaskID, googlecloudclustercomposer_contract.AirflowWorkerLogFilterTaskID.Ref(), )
var AirflowWorkerLogToTimelineMapperTask = inspectiontaskbase.NewLogToTimelineMapperTask[struct{}]( googlecloudclustercomposer_contract.AirflowWorkerLogToTimelineMapperTaskID, &airflowWorkerLogToTimelineMapperSetting{ targetLogType: enum.LogTypeComposerEnvironment, }, )
var AutocompleteComposerClusterNamesTask = inspectiontaskbase.NewCachedTask(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(), }, 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()) dependencyDigest := fmt.Sprintf("%s-%s-%s", projectID, environment, location) 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 environment != "" && dependencyDigest == prevValue.DependencyDigest { return prevValue, nil } clusterFinder := coretask.GetTaskResult(ctx, googlecloudclustercomposer_contract.ComposerEnvironmentClusterFinderTaskID.Ref()) clusterName, err := clusterFinder.GetGKEClusterName(ctx, projectID, environment) 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 } return inspectiontaskbase.CacheableTaskResult[*inspectioncore_contract.AutocompleteResult[googlecloudk8scommon_contract.GoogleCloudClusterIdentity]]{ DependencyDigest: dependencyDigest, Value: &inspectioncore_contract.AutocompleteResult[googlecloudk8scommon_contract.GoogleCloudClusterIdentity]{ Values: []googlecloudk8scommon_contract.GoogleCloudClusterIdentity{ { ClusterName: clusterName, ProjectID: projectID, Location: location, }, }, }, }, 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.
var AutocompleteComposerComponentsTask = inspectiontaskbase.NewCachedTask(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) } defer client.Close() 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 })
var AutocompleteComposerEnvironmentIdentityTask = inspectiontaskbase.NewCachedTask(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) } defer client.Close() 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.
var AutocompleteLocationForComposerEnvironmentTask = inspectiontaskbase.NewCachedTask(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), )
var ClusterIdentityAliasTask = coretask.NewAliasTask( googlecloudclustercomposer_contract.ClusterIdentityTaskID, googlecloudk8scommon_contract.ClusterIdentityTaskID.Ref(), )
var ComposerEnvironmentClusterFinderTask = coretask.NewTask( googlecloudclustercomposer_contract.ComposerEnvironmentClusterFinderTaskID, []taskid.UntypedTaskReference{ googlecloudcommon_contract.APIClientFactoryTaskID.Ref(), }, func(ctx context.Context) (googlecloudclustercomposer_contract.ComposerEnvironmentClusterFinder, error) { return &googlecloudclustercomposer_contract.EnvironmentClusterFinderImpl{}, nil }, )
ComposerEnvironmentClusterFinderTask injects ComposerEnvironmentClusterFinder implementation.
var ComposerEnvironmentListFetcherTask = coretask.NewTask( googlecloudclustercomposer_contract.ComposerEnvironmentListFetcherTaskID, []taskid.UntypedTaskReference{ googlecloudcommon_contract.APIClientFactoryTaskID.Ref(), }, func(ctx context.Context) (googlecloudclustercomposer_contract.ComposerEnvironmentListFetcher, error) { return &googlecloudclustercomposer_contract.ComposerEnvironmentListFetcherImpl{}, nil }, )
ComposerEnvironmentListFetcherTask injects ComposerEnvironmentListFetcher implementation.
var ComposerLogsFieldSetReadTask = inspectiontaskbase.NewFieldSetReadTask( googlecloudclustercomposer_contract.ComposerLogsFieldSetReadTaskID, googlecloudclustercomposer_contract.ComposerLogsQueryTaskID.Ref(), []log.FieldSetReader{ &gcpqueryutil.GCPMainMessageFieldSetReader{}, &googlecloudclustercomposer_contract.ComposerFieldSetReader{}, &googlecloudclustercomposer_contract.ComposerTaskInstanceFieldSetReader{}, &googlecloudclustercomposer_contract.ComposerWorkerTaskInstanceFieldSetReader{}, }, )
ComposerLogsFieldSetReadTask reads the main message and Composer component fieldsets.
var ComposerLogsQueryTask = googlecloudcommon_contract.NewListLogEntriesTask(&composerListLogEntriesTaskSetting{ taskId: googlecloudclustercomposer_contract.ComposerLogsQueryTaskID, queryName: "Composer Environment Logs", })
ComposerLogsQueryTask defines a task that gathers logs from Cloud Logging for multiple Composer components.
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( "Composer Logs", "Cloud Composer related logs like airflow-worker, airflow-scheduler, airflow-dag-processor-manager, and others.", enum.LogTypeComposerEnvironment, 101000, true, ), )
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()
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 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.
Source Files
¶
- airflow_util.go
- autocomplete.go
- autocompletecomposerclusternames_task.go
- clusteridentity_task.go
- dag_processor_manager_mapper.go
- fieldsetreader_task.go
- filter_task.go
- inject_task.go
- input_components_task.go
- inputcomposerenvironmentname_task.go
- other.go
- query_task.go
- registration.go
- scheduler.go
- tail_task.go
- worker.go