Documentation
¶
Index ¶
- Variables
- func GenerateComposerLogsQuery(projectID, location, environmentName string, selectedComponents []string) string
- func GenerateComposerLogsStructuredQuery(projectID, location, environmentName string, selectedComponents []string) *logestimator.StructuredLogQuery
- func Register(registry coreinspection.InspectionTaskRegistry) error
- type DagProcessorState
Constants ¶
This section is empty.
Variables ¶
var AirflowDagProcessorManagerLogFilterTask = componentFilterTask(googlecloudclustercomposer_contract.AirflowDagProcessorManagerLogFilterTaskID, googlecloudclustercomposer_contract.ComposerLogsQueryTaskID.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 := 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.
var AirflowDagProcessorManagerLogIngesterTask = inspectiontaskbase.NewGroupedLogIngesterTask( googlecloudclustercomposer_contract.AirflowDagProcessorManagerLogIngesterTaskID, &dagProcessorManagerLogIngester{}, )
AirflowDagProcessorManagerLogIngesterTask is the task that ingests Airflow DAG processor manager logs.
var AirflowDagProcessorManagerLogSorterTask = inspectiontaskbase.NewLogSorterByTimeTask( googlecloudclustercomposer_contract.AirflowDagProcessorManagerLogSorterTaskID, googlecloudclustercomposer_contract.AirflowDagProcessorManagerLogFilterTaskID.Ref(), )
AirflowDagProcessorManagerLogSorterTask sorts Airflow DAG processor manager logs.
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.
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" }, )
var AirflowOtherLogGrouperTask = inspectiontaskbase.NewLogGrouperTask( googlecloudclustercomposer_contract.AirflowOtherLogGrouperTaskID, googlecloudclustercomposer_contract.AirflowOtherLogFilterTaskID.Ref(), func(ctx context.Context, l *log.Log) string { return "" }, )
AirflowOtherLogGrouperTask groups other Airflow logs.
var AirflowOtherLogIngesterTask = inspectiontaskbase.NewLogIngesterTask( googlecloudclustercomposer_contract.AirflowOtherLogIngesterTaskID, &otherLogIngester{}, )
AirflowOtherLogIngesterTask is the task that ingests other Airflow logs.
var AirflowOtherLogToTimelineMapperTask = inspectiontaskbase.NewLogToTimelineMapperTask( googlecloudclustercomposer_contract.AirflowOtherLogToTimelineMapperTaskID, &otherLogToTimelineMapper{}, )
AirflowOtherLogToTimelineMapperTask is the task that maps other Airflow logs to timeline events.
var AirflowSchedulerLogFilterTask = componentFilterTask(googlecloudclustercomposer_contract.AirflowSchedulerLogFilterTaskID, googlecloudclustercomposer_contract.ComposerLogsQueryTaskID.Ref(), "airflow-scheduler")
var AirflowSchedulerLogGrouperTask = inspectiontaskbase.NewLogGrouperTask( googlecloudclustercomposer_contract.AirflowSchedulerLogGrouperTaskID, googlecloudclustercomposer_contract.AirflowSchedulerLogFilterTaskID.Ref(), func(ctx context.Context, l *log.Log) string { return "" }, )
AirflowSchedulerLogGrouperTask groups Airflow scheduler logs.
var AirflowSchedulerLogIngesterTask = inspectiontaskbase.NewLogIngesterTask( googlecloudclustercomposer_contract.AirflowSchedulerLogIngesterTaskID, &schedulerLogIngester{}, )
AirflowSchedulerLogIngesterTask is the task that ingests Airflow scheduler logs.
var AirflowSchedulerLogToTimelineMapperTask = inspectiontaskbase.NewLogToTimelineMapperTask( googlecloudclustercomposer_contract.AirflowSchedulerLogToTimelineMapperTaskID, &schedulerLogToTimelineMapper{}, )
AirflowSchedulerLogToTimelineMapperTask is the task that maps Airflow scheduler logs to timeline events.
var AirflowWorkerLogFilterTask = componentFilterTask(googlecloudclustercomposer_contract.AirflowWorkerLogFilterTaskID, googlecloudclustercomposer_contract.ComposerLogsQueryTaskID.Ref(), "airflow-worker")
var AirflowWorkerLogGrouperTask = inspectiontaskbase.NewLogGrouperTask( googlecloudclustercomposer_contract.AirflowWorkerLogGrouperTaskID, googlecloudclustercomposer_contract.AirflowWorkerLogFilterTaskID.Ref(), func(ctx context.Context, l *log.Log) string { return "" }, )
AirflowWorkerLogGrouperTask groups Airflow worker logs.
var AirflowWorkerLogIngesterTask = inspectiontaskbase.NewLogIngesterTask( googlecloudclustercomposer_contract.AirflowWorkerLogIngesterTaskID, &workerLogIngester{}, )
AirflowWorkerLogIngesterTask is the task that ingests Airflow worker logs.
var AirflowWorkerLogToTimelineMapperTask = inspectiontaskbase.NewLogToTimelineMapperTask( googlecloudclustercomposer_contract.AirflowWorkerLogToTimelineMapperTaskID, &workerLogToTimelineMapper{}, )
AirflowWorkerLogToTimelineMapperTask is the task that maps Airflow worker logs to timeline events.
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.
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 })
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.
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), )
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 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.
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, ), )
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 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.