googlecloudclustercomposer_contract

package
v0.53.0 Latest Latest
Warning

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

Go to latest
Published: Apr 21, 2026 License: Apache-2.0 Imports: 21 Imported by: 0

Documentation

Index

Constants

View Source
const InspectionTypeID = "gcp-composer"

InspectionTypeID is the inspection type id for google cloud composer.

Variables

View Source
var AirflowDagProcessorManagerLogFilterTaskID taskid.TaskImplementationID[[]*log.Log] = taskid.NewDefaultImplementationID[[]*log.Log](GoogleCloudComposerTaskIDPrefix + "filter-dag-processor-manager")

AirflowDagProcessorManagerLogFilterTaskID is the task id for filtering Airflow DAG processor manager logs.

View Source
var AirflowDagProcessorManagerLogGrouperTaskID = taskid.NewDefaultImplementationID[inspectiontaskbase.LogGroupMap](GoogleCloudComposerTaskIDPrefix + "grouper-dag-processor-manager")

AirflowDagProcessorManagerLogGrouperTaskID is the task id for the task that groups Airflow DAG processor manager logs.

View Source
var AirflowDagProcessorManagerLogIngesterTaskID = taskid.NewDefaultImplementationID[[]*log.Log](GoogleCloudComposerTaskIDPrefix + "ingester-dag-processor-manager")

AirflowDagProcessorManagerLogIngesterTaskID is the task id for the task that ingests Airflow DAG processor manager logs.

View Source
var AirflowDagProcessorManagerLogSorterTaskID = taskid.NewDefaultImplementationID[[]*log.Log](GoogleCloudComposerTaskIDPrefix + "sorter-dag-processor-manager")

AirflowDagProcessorManagerLogSorterTaskID is the task id for the task that sorts Airflow DAG processor manager logs.

View Source
var AirflowDagProcessorManagerLogToTimelineMapperTaskID = taskid.NewDefaultImplementationID[struct{}](GoogleCloudComposerTaskIDPrefix + "mapper-dag-processor-manager")

AirflowDagProcessorManagerLogToTimelineMapperTaskID is the task id for the task that maps Airflow DAG processor manager logs to timeline events.

AirflowOtherLogFilterTaskID is the task id for filtering other Airflow logs.

AirflowOtherLogGrouperTaskID is the task id for the task that groups other Airflow logs.

View Source
var AirflowOtherLogIngesterTaskID = taskid.NewDefaultImplementationID[[]*log.Log](GoogleCloudComposerTaskIDPrefix + "ingester-other")

AirflowOtherLogIngesterTaskID is the task id for the task that ingests other Airflow logs.

View Source
var AirflowOtherLogToTimelineMapperTaskID = taskid.NewDefaultImplementationID[struct{}](GoogleCloudComposerTaskIDPrefix + "mapper-other")

AirflowOtherLogToTimelineMapperTaskID is the task id for the task that maps other Airflow logs to timeline events.

View Source
var AirflowSchedulerLogFilterTaskID taskid.TaskImplementationID[[]*log.Log] = taskid.NewDefaultImplementationID[[]*log.Log](GoogleCloudComposerTaskIDPrefix + "filter-scheduler")

AirflowSchedulerLogFilterTaskID is the task id for filtering Airflow scheduler logs.

View Source
var AirflowSchedulerLogGrouperTaskID = taskid.NewDefaultImplementationID[inspectiontaskbase.LogGroupMap](GoogleCloudComposerTaskIDPrefix + "grouper-scheduler")

AirflowSchedulerLogGrouperTaskID is the task id for the task that groups Airflow scheduler logs.

View Source
var AirflowSchedulerLogIngesterTaskID = taskid.NewDefaultImplementationID[[]*log.Log](GoogleCloudComposerTaskIDPrefix + "ingester-scheduler")

AirflowSchedulerLogIngesterTaskID is the task id for the task that ingests Airflow scheduler logs.

View Source
var AirflowSchedulerLogToTimelineMapperTaskID = taskid.NewDefaultImplementationID[struct{}](GoogleCloudComposerTaskIDPrefix + "mapper-scheduler")

AirflowSchedulerLogToTimelineMapperTaskID is the task id for the task that maps Airflow scheduler logs to timeline events.

AirflowWorkerLogFilterTaskID is the task id for filtering Airflow worker logs.

AirflowWorkerLogGrouperTaskID is the task id for the task that groups Airflow worker logs.

View Source
var AirflowWorkerLogIngesterTaskID = taskid.NewDefaultImplementationID[[]*log.Log](GoogleCloudComposerTaskIDPrefix + "ingester-worker")

AirflowWorkerLogIngesterTaskID is the task id for the task that ingests Airflow worker logs.

View Source
var AirflowWorkerLogToTimelineMapperTaskID = taskid.NewDefaultImplementationID[struct{}](GoogleCloudComposerTaskIDPrefix + "mapper-worker")

AirflowWorkerLogToTimelineMapperTaskID is the task id for the task that maps Airflow worker logs to timeline events.

View Source
var AutocompleteComposerClusterNamesTaskID = taskid.NewImplementationID(googlecloudk8scommon_contract.AutocompleteClusterIdentityTaskID.Ref(), "composer")

AutocompleteComposerClusterNamesTaskID is the task id for the task that autocompletes GKE cluster names created by Cloud Composer.

View Source
var AutocompleteComposerComponentsTaskID = taskid.NewDefaultImplementationID[*inspectioncore_contract.AutocompleteResult[string]](GoogleCloudComposerTaskIDPrefix + "autocomplete/composer-components")

AutocompleteComposerComponentsTaskID is the task id for autocompleting component names from Cloud Monitoring.

View Source
var AutocompleteComposerEnvironmentIdentityTaskID = taskid.NewDefaultImplementationID[*inspectioncore_contract.AutocompleteResult[ComposerEnvironmentIdentity]](GoogleCloudComposerTaskIDPrefix + "autocomplete/composer-environment-identities")

AutocompleteComposerEnvironmentIdentityTaskID is the task id for the task that autocompletes composer environment identities.

View Source
var AutocompleteComposerEnvironmentNamesTaskID taskid.TaskImplementationID[[]string] = taskid.NewDefaultImplementationID[[]string](GoogleCloudComposerTaskIDPrefix + "autocomplete/composer-environment-names")

AutocompleteComposerEnvironmentNamesTaskID is the task id for the task that autocompletes composer environment names.

View Source
var AutocompleteLocationForComposerEnvironmentTaskID = taskid.NewImplementationID(googlecloudcommon_contract.AutocompleteLocationTaskID.Ref(), "composer")

AutocompleteLocationForComposerEnvironmentTaskID is the task id for the task that autocompletes GKE cluster location from Composer environments.

ClusterIdentityTaskID is the task id for aliasing the cluster identity.

ComposerClusterNamePrefixTaskID is the task id for the task that returns the GKE cluster name prefix used by Cloud Composer.

View Source
var ComposerEnvironmentClusterFinderTaskID = taskid.NewDefaultImplementationID[ComposerEnvironmentClusterFinder](GoogleCloudComposerTaskIDPrefix + "composer-environment-cluster-finder")

ComposerEnvironmentClusterFinderTaskID is the task id for injecting ComposerEnvironmentClusterFinder instance.

View Source
var ComposerEnvironmentListFetcherTaskID = taskid.NewDefaultImplementationID[ComposerEnvironmentListFetcher](GoogleCloudComposerTaskIDPrefix + "composer-environment-list-fetcher")

ComposerEnvironmentListFetcherTaskID is the task id for injecting ComposerEnvironmentListFetcher instance.

View Source
var ComposerInspectionType = coreinspection.InspectionType{
	Id:   InspectionTypeID,
	Name: "Cloud Composer",
	Description: `Visualize logs related to Cloud Composer environment.
Supports all GKE related logs(Cloud Composer v2) and Airflow logs(Airflow 2.0.0 or higher in any Cloud Composer version(v1-v2, partical v3))`,
	Icon:     "assets/icons/composer.webp",
	Priority: math.MaxInt - 10,
	Labels: map[string]string{
		inspectioncore_contract.InspectionTypeLabelKeyLogSource:      "cloud_logging",
		inspectioncore_contract.InspectionTypeLabelKeyEnvironment:    "googlecloud",
		inspectioncore_contract.InspectionTypeLabelKeyBasePlatform:   "kubernetes",
		googlecloudcommon_contract.InspectionTypeLabelKeyClusterType: "gke",
		googlecloudcommon_contract.InspectionTypeLabelKeyProduct:     "composer",
	},
}

ComposerInspectionType is the inspection type for google cloud composer.

ComposerLogsFieldSetReadTaskID is the task id for the task that reads fieldsets from composer logs.

ComposerLogsQueryTaskID is the task id for the task that queries Logs from Cloud Logging.

View Source
var ComposerLogsTailTaskID = taskid.NewDefaultImplementationID[struct{}](GoogleCloudComposerTaskIDPrefix + "tail-composer-logs")

ComposerLogsTailTaskID is the task id for unifying composer logs feature.

View Source
var ErrEnvironmentClusterNotFound = errors.New("not found")
View Source
var GoogleCloudComposerTaskIDPrefix = "cloud.google.com/composer/"

GoogleCloudComposerTaskIDPrefix is the prefix for all task ids related to google cloud composer.

View Source
var InputComposerComponentsTaskID taskid.TaskImplementationID[[]string] = taskid.NewDefaultImplementationID[[]string](GoogleCloudComposerTaskIDPrefix + "input/composer/components")

InputComposerComponentsTaskID is the task id for selecting target Composer components.

View Source
var InputComposerEnvironmentNameTaskID taskid.TaskImplementationID[string] = taskid.NewDefaultImplementationID[string](GoogleCloudComposerTaskIDPrefix + "input/composer/environment_name")

InputComposerEnvironmentNameTaskID is the task id for the task that inputs composer environment name.

Functions

This section is empty.

Types

type AirflowTaskInstance added in v0.52.8

type AirflowTaskInstance struct {
	// contains filtered or unexported fields
}

ref: https://github.com/apache/airflow/blob/main/airflow/models/taskinstance.py#L1187

func NewAirflowTaskInstance added in v0.52.8

func NewAirflowTaskInstance(dagId string, taskId string, runId string, mapIndex string, host string, status Tistate) *AirflowTaskInstance

func (*AirflowTaskInstance) DagId added in v0.52.8

func (a *AirflowTaskInstance) DagId() string

func (*AirflowTaskInstance) Host added in v0.52.8

func (a *AirflowTaskInstance) Host() string

func (*AirflowTaskInstance) MapIndex added in v0.52.8

func (a *AirflowTaskInstance) MapIndex() string

func (*AirflowTaskInstance) ResourcePath added in v0.52.8

func (a *AirflowTaskInstance) ResourcePath() resourcepath.ResourcePath

func (*AirflowTaskInstance) RunId added in v0.52.8

func (a *AirflowTaskInstance) RunId() string

func (*AirflowTaskInstance) Status added in v0.52.8

func (a *AirflowTaskInstance) Status() Tistate

func (*AirflowTaskInstance) TaskId added in v0.52.8

func (a *AirflowTaskInstance) TaskId() string

func (*AirflowTaskInstance) ToYaml added in v0.52.8

func (a *AirflowTaskInstance) ToYaml() string

type AirflowWorker added in v0.52.8

type AirflowWorker struct {
	// contains filtered or unexported fields
}

func NewAirflowWorker added in v0.52.8

func NewAirflowWorker(host string) *AirflowWorker

func (*AirflowWorker) Host added in v0.52.8

func (a *AirflowWorker) Host() string

func (*AirflowWorker) ResourcePath added in v0.52.8

func (a *AirflowWorker) ResourcePath() resourcepath.ResourcePath

func (*AirflowWorker) ToYaml added in v0.52.8

func (a *AirflowWorker) ToYaml() string

type ComposerEnvironmentClusterFinder

type ComposerEnvironmentClusterFinder interface {
	GetGKEClusterName(ctx context.Context, projectID, environment string) (string, error)
}

type ComposerEnvironmentIdentity added in v0.52.0

type ComposerEnvironmentIdentity struct {
	ProjectID       string
	Location        string
	EnvironmentName string
}

ComposerEnvironmentIdentity represents the identity of a Cloud Composer environment.

type ComposerEnvironmentListFetcher

type ComposerEnvironmentListFetcher interface {
	GetEnvironmentNames(ctx context.Context, projectID, location string) ([]string, error)
}

ComposerEnvironmentListFetcher fetches the list of Cloud Composer environment names in project and the location.

type ComposerEnvironmentListFetcherImpl

type ComposerEnvironmentListFetcherImpl struct{}

func (*ComposerEnvironmentListFetcherImpl) GetEnvironmentNames

func (c *ComposerEnvironmentListFetcherImpl) GetEnvironmentNames(ctx context.Context, projectID string, location string) ([]string, error)

GetEnvironmentNames implements ComposerEnvironmentListFetcher.

type ComposerFieldSet added in v0.52.8

type ComposerFieldSet struct {
	Component             string // e.g. "worker", "scheduler", "dag-processor-manager"
	WorkerID              string
	SchedulerID           string
	DagProcessorManagerID string
	TriggererID           string
	WebserverID           string
	Subservice            string
}

func (*ComposerFieldSet) Kind added in v0.52.8

func (c *ComposerFieldSet) Kind() string

type ComposerFieldSetReader added in v0.52.8

type ComposerFieldSetReader struct{}

func (*ComposerFieldSetReader) FieldSetKind added in v0.52.8

func (c *ComposerFieldSetReader) FieldSetKind() string

func (*ComposerFieldSetReader) Read added in v0.52.8

type ComposerTaskInstanceFieldSet added in v0.52.8

type ComposerTaskInstanceFieldSet struct {
	TaskInstance *AirflowTaskInstance
}

func (*ComposerTaskInstanceFieldSet) Kind added in v0.52.8

type ComposerTaskInstanceFieldSetReader added in v0.52.8

type ComposerTaskInstanceFieldSetReader struct{}

func (*ComposerTaskInstanceFieldSetReader) FieldSetKind added in v0.52.8

func (c *ComposerTaskInstanceFieldSetReader) FieldSetKind() string

func (*ComposerTaskInstanceFieldSetReader) Read added in v0.52.8

type ComposerWorkerTaskInstanceFieldSet added in v0.52.8

type ComposerWorkerTaskInstanceFieldSet struct {
	TaskInstance *AirflowTaskInstance
}

func (*ComposerWorkerTaskInstanceFieldSet) Kind added in v0.52.8

type ComposerWorkerTaskInstanceFieldSetReader added in v0.52.8

type ComposerWorkerTaskInstanceFieldSetReader struct{}

func (*ComposerWorkerTaskInstanceFieldSetReader) FieldSetKind added in v0.52.8

func (*ComposerWorkerTaskInstanceFieldSetReader) Read added in v0.52.8

type EnvironmentClusterFinderImpl

type EnvironmentClusterFinderImpl struct{}

func (*EnvironmentClusterFinderImpl) GetGKEClusterName

func (e *EnvironmentClusterFinderImpl) GetGKEClusterName(ctx context.Context, projectID string, environment string) (string, error)

GetGKEClusterName implements EnvironmentClusterFinder.

type Tistate added in v0.52.8

type Tistate string
const (
	// ref: https://airflow.apache.org/docs/apache-airflow/stable/core-concepts/tasks.html#task-instances
	TASKINSTANCE_NONE              Tistate = "none"
	TASKINSTANCE_SCHEDULED         Tistate = "scheduled"
	TASKINSTANCE_QUEUED            Tistate = "queued"
	TASKINSTANCE_RUNNING           Tistate = "running"
	TASKINSTANCE_SUCCESS           Tistate = "success"
	TASKINSTANCE_SHUTDOWN          Tistate = "shutdown"
	TASKINSTANCE_RESTARTING        Tistate = "restarting"
	TASKINSTANCE_FAILED            Tistate = "failed"
	TASKINSTANCE_SKIPPED           Tistate = "skipped"
	TASKINSTANCE_UP_FOR_RETRY      Tistate = "up_for_retry"
	TASKINSTANCE_DEFERRED          Tistate = "deferred"
	TASKINSTANCE_UP_FOR_RESCHEDULE Tistate = "up_for_reschedule"
	TASKINSTANCE_REMOVED           Tistate = "removed"
	TASKINSTANCE_UPSTREAM_FAILED   Tistate = "upstream_failed"

	// Original States //
	// Zombie status for KHI view
	TASKINSTANCE_ZOMBIE Tistate = "zombie"
)

Jump to

Keyboard shortcuts

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