googlecloudclustercomposer_contract

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: 25 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 (
	RevisionStateComposerTiScheduled = style.MustRegisterRevisionState(
		"Task instance is scheduled",
		"schedule",
		"The Airflow task instance has been scheduled and is waiting to be queued.",
		style.MustForceConvertSRGBHex("#d1b48c"),
		pb.RevisionStateStyle_REVISION_STATE_STYLE_NORMAL,
	)
	RevisionStateComposerTiQueued = style.MustRegisterRevisionState(
		"Task instance is queued",
		"transition_push",
		"The Airflow task instance has been queued in the executor and is waiting to run.",
		style.MustForceConvertSRGBHex("#808080"),
		pb.RevisionStateStyle_REVISION_STATE_STYLE_NORMAL,
	)
	RevisionStateComposerTiRunning = style.MustRegisterRevisionState(
		"Task instance is running",
		"directions_run",
		"The Airflow task instance is currently executing.",
		style.MustForceConvertSRGBHex("#00ff01"),
		pb.RevisionStateStyle_REVISION_STATE_STYLE_NORMAL,
	)
	RevisionStateComposerTiDeferred = style.MustRegisterRevisionState(
		"Task instance is deferred",
		"pause",
		"The Airflow task instance is deferred, waiting for a trigger to resume.",
		style.MustForceConvertSRGBHex("#9470dc"),
		pb.RevisionStateStyle_REVISION_STATE_STYLE_NORMAL,
	)
	RevisionStateComposerTiSuccess = style.MustRegisterRevisionState(
		"Task instance succeeded",
		"check",
		"The Airflow task instance has completed successfully.",
		style.MustForceConvertSRGBHex("#008001"),
		pb.RevisionStateStyle_REVISION_STATE_STYLE_NORMAL,
	)
	RevisionStateComposerTiFailed = style.MustRegisterRevisionState(
		"Task instance failed",
		"exclamation",
		"The Airflow task instance has failed during execution.",
		style.MustForceConvertSRGBHex("#fe0000"),
		pb.RevisionStateStyle_REVISION_STATE_STYLE_NORMAL,
	)
	RevisionStateComposerTiUpForRetry = style.MustRegisterRevisionState(
		"Task instance is up for retry",
		"camping",
		"The Airflow task instance has failed and is waiting to be retried.",
		style.MustForceConvertSRGBHex("#fed700"),
		pb.RevisionStateStyle_REVISION_STATE_STYLE_NORMAL,
	)
	RevisionStateComposerTiRestarting = style.MustRegisterRevisionState(
		"Task instance is restarting",
		"restart_alt",
		"The Airflow task instance is being restarted.",
		style.MustForceConvertSRGBHex("#ee82ef"),
		pb.RevisionStateStyle_REVISION_STATE_STYLE_NORMAL,
	)
	RevisionStateComposerTiRemoved = style.MustRegisterRevisionState(
		"Task instance is removed",
		"waving_hand",
		"The Airflow task instance has been removed from the DAG run.",
		style.MustForceConvertSRGBHex("#d3d3d3"),
		pb.RevisionStateStyle_REVISION_STATE_STYLE_NORMAL,
	)
	RevisionStateComposerTiUpstreamFailed = style.MustRegisterRevisionState(
		"Upstream task failed",
		"falling",
		"The Airflow task instance has been skipped because one of its upstream dependencies failed.",
		style.MustForceConvertSRGBHex("#ffa11b"),
		pb.RevisionStateStyle_REVISION_STATE_STYLE_NORMAL,
	)
	RevisionStateComposerTiZombie = style.MustRegisterRevisionState(
		"Task instance is a zombie",
		"skull",
		"The Airflow task instance is detected as a zombie (the process died without updating the database state).",
		style.MustForceConvertSRGBHex("#4b0082"),
		pb.RevisionStateStyle_REVISION_STATE_STYLE_NORMAL,
	)
	RevisionStateComposerTiUpForReschedule = style.MustRegisterRevisionState(
		"Task instance is up for reschedule",
		"history",
		"The Airflow task instance is in up_for_reschedule state, waiting for the next sensor poll.",
		style.MustForceConvertSRGBHex("#808080"),
		pb.RevisionStateStyle_REVISION_STATE_STYLE_NORMAL,
	)
	RevisionStateComposerTiSkipped = style.MustRegisterRevisionState(
		"Task instance is skipped",
		"step_over",
		"The Airflow task instance has been skipped during execution.",
		style.MustForceConvertSRGBHex("#e60076"),
		pb.RevisionStateStyle_REVISION_STATE_STYLE_NORMAL,
	)
	// RevisionStateComposerDagProcessorNoError indicates that DAG file processing completed without errors.
	RevisionStateComposerDagProcessorNoError = style.MustRegisterRevisionState(
		"DAG processing has no errors",
		"check",
		"The Airflow DAG processor manager processed the DAG file without errors.",
		style.MustForceConvertSRGBHex("#008001"),
		pb.RevisionStateStyle_REVISION_STATE_STYLE_NORMAL,
	)
	// RevisionStateComposerDagProcessorHasErrors indicates that DAG file processing encountered errors.
	RevisionStateComposerDagProcessorHasErrors = style.MustRegisterRevisionState(
		"DAG processing has errors",
		"exclamation",
		"The Airflow DAG processor manager encountered errors while processing the DAG file.",
		style.MustForceConvertSRGBHex("#fe0000"),
		pb.RevisionStateStyle_REVISION_STATE_STYLE_NORMAL,
	)
	// RevisionStateManagedAirflowEnvironmentProvisioning indicates that the Managed Airflow environment is being created.
	RevisionStateManagedAirflowEnvironmentProvisioning = style.MustRegisterRevisionState(
		"Environment is being provisioned",
		"deployed_code_history",
		"The Managed Airflow environment is currently being provisioned.",
		style.MustForceConvertSRGBHex("#6666ff"),
		pb.RevisionStateStyle_REVISION_STATE_STYLE_NORMAL,
	)
	// RevisionStateManagedAirflowEnvironmentExisting indicates that the Managed Airflow environment is active and running.
	RevisionStateManagedAirflowEnvironmentExisting = style.MustRegisterRevisionState(
		"Environment exists",
		"deployed_code",
		"The Managed Airflow environment exists and is active.",
		style.Color{R: 0.0, G: 0.0, B: 1.0, A: 1.0},
		pb.RevisionStateStyle_REVISION_STATE_STYLE_NORMAL,
	)
	// RevisionStateManagedAirflowEnvironmentDeleting indicates that the Managed Airflow environment is being deleted.
	RevisionStateManagedAirflowEnvironmentDeleting = style.MustRegisterRevisionState(
		"Environment is being deleted",
		"auto_delete",
		"The Managed Airflow environment is undergoing deletion.",
		style.Color{R: 0.8, G: 0.33333334, B: 0.0, A: 1.0},
		pb.RevisionStateStyle_REVISION_STATE_STYLE_NORMAL,
	)
	// RevisionStateManagedAirflowEnvironmentDeleted indicates that the Managed Airflow environment has been deleted.
	RevisionStateManagedAirflowEnvironmentDeleted = style.MustRegisterRevisionState(
		"Environment is deleted",
		"delete_forever",
		"The Managed Airflow environment has been deleted.",
		style.Color{R: 0.8, G: 0.0, B: 0.0, A: 1.0},
		pb.RevisionStateStyle_REVISION_STATE_STYLE_DELETED,
	)
	// RevisionStateManagedAirflowEnvironmentProvisioningLogNotFound indicates that Managed Airflow environment provisioning started before the log collection window.
	RevisionStateManagedAirflowEnvironmentProvisioningLogNotFound = style.MustRegisterRevisionState(
		"Environment is being provisioned, but starting log not found",
		"deployed_code_history",
		"The Managed Airflow environment provisioning was started, but the starting log entry was not found in the selected time range.",
		style.MustForceConvertSRGBHex("#6666ff"),
		pb.RevisionStateStyle_REVISION_STATE_STYLE_PARTIAL_INFO,
	)
	// RevisionStateManagedAirflowEnvironmentExistingLogNotFound indicates that Managed Airflow environment existed before the log collection window.
	RevisionStateManagedAirflowEnvironmentExistingLogNotFound = style.MustRegisterRevisionState(
		"Environment exists, but creation log not found",
		"deployed_code",
		"The Managed Airflow environment exists, but the creation or existence log entry was not found in the selected time range.",
		style.Color{R: 0.0, G: 0.0, B: 1.0, A: 1.0},
		pb.RevisionStateStyle_REVISION_STATE_STYLE_PARTIAL_INFO,
	)
	// RevisionManagedAirflowEnvironmentDeletingLogNotFound indicates that Managed Airflow environment deletion started before the log collection window.
	RevisionManagedAirflowEnvironmentDeletingLogNotFound = style.MustRegisterRevisionState(
		"Environment is being deleted, but starting log not found",
		"auto_delete",
		"The Managed Airflow environment deletion was in progress, but the deletion starting log entry was not found in the selected time range.",
		style.Color{R: 0.8, G: 0.33333334, B: 0.0, A: 1.0},
		pb.RevisionStateStyle_REVISION_STATE_STYLE_PARTIAL_INFO,
	)
)

The following block defines the registered timeline style RevisionStates. These are registered as package-level variables so they are initialized immediately when this package is imported.

View Source
var (

	// TimelineTypeAirflow is the style for an Airflow environment.
	TimelineTypeAirflow = style.MustRegisterTimelineType(
		"Airflow",
		"Timeline representing a Managed Airflow environment",
		"settings",
		0.6,
		style.Color{R: 0.82, G: 0.28, B: 0.19, A: 1},
		style.ColorWhite,
		style.MustForceConvertSRGBHex("#F5F5F5"),
		style.ColorBlack,
		true,
		20,
		style.AlphabeticalSortPolicy(),
	)

	// TimelineTypeDAGs is the root style for DAGs hierarchy.
	TimelineTypeDAGs = style.MustRegisterTimelineType(
		"DAGs",
		"Grouping timeline for Airflow DAGs",
		"folder",
		0.6,
		style.Color{R: 0.93, G: 0.49, B: 0.38, A: 1},
		style.ColorWhite,
		style.MustForceConvertSRGBHex("#F5F5F5"),
		style.ColorBlack,
		true,
		30,
		style.AlphabeticalSortPolicy(),
	)

	// TimelineTypeAirflowDAG is the style for a single DAG.
	TimelineTypeAirflowDAG = style.MustRegisterTimelineType(
		"DAG",
		"Timeline representing an Airflow DAG",
		"account_tree",
		0.6,
		style.MustForceConvertSRGBHex("#444444"),
		style.ColorWhite,
		style.MustForceConvertSRGBHex("#F5F5F5"),
		style.ColorBlack,
		true,
		40,
		style.ChronologicalSortPolicy(2),
	)

	// TimelineTypeAirflowDAGRun is the style for a DAG run.
	TimelineTypeAirflowDAGRun = style.MustRegisterTimelineType(
		"DAG Run",
		"Timeline representing an Airflow DAG run",
		"play_circle",
		0.6,
		style.MustForceConvertSRGBHex("#CCCCCC"),
		style.ColorBlack,
		style.MustForceConvertSRGBHex("#F5F5F5"),
		style.ColorBlack,
		true,
		50,
		style.ChronologicalSortPolicy(1),
	)

	// TimelineTypeAirflowTaskInstance is the style for a TaskInstance.
	TimelineTypeAirflowTaskInstance = style.MustRegisterTimelineType(
		"Task Instance",
		"Execution states of the Airflow task instance",
		"mode_fan",
		0.7,
		style.ColorWhite,
		style.ColorBlack,
		style.MustForceConvertSRGBHex("#F5F5F5"),
		style.ColorBlack,
		true,
		1501,
		style.ChronologicalSortPolicy(1),
	)

	// TimelineTypeComponents is the category style for components.
	TimelineTypeComponents = style.MustRegisterTimelineType(
		"Components",
		"Grouping timeline for Airflow backend components",
		"apps",
		0.6,
		style.Color{R: 0.93, G: 0.49, B: 0.38, A: 1},
		style.ColorWhite,
		style.MustForceConvertSRGBHex("#F5F5F5"),
		style.ColorBlack,
		true,
		70,
		style.AlphabeticalSortPolicy(),
	)

	// TimelineTypeDAGFiles is the category style for DAG Processor Manager stats.
	TimelineTypeDAGFiles = style.MustRegisterTimelineType(
		"DAG files",
		"Grouping timeline for parsed DAG files",
		"folder",
		0.6,
		style.Color{R: 0.93, G: 0.49, B: 0.38, A: 1},
		style.ColorWhite,
		style.MustForceConvertSRGBHex("#F5F5F5"),
		style.ColorBlack,
		true,
		60,
		style.AlphabeticalSortPolicy(),
	)

	// TimelineTypeDAGFile is the style for a parsed DAG file.
	TimelineTypeDAGFile = style.MustRegisterTimelineType(
		"DAG File",
		"Timeline representing an Airflow DAG definition file",
		"description",
		0.6,
		style.MustForceConvertSRGBHex("#444444"),
		style.ColorWhite,
		style.MustForceConvertSRGBHex("#F5F5F5"),
		style.ColorBlack,
		true,
		80,
		style.AlphabeticalSortPolicy(),
	)

	// TimelineTypeDAGProcessorManagerInstance is the style for the manager instance that processed the file.
	TimelineTypeDAGProcessorManagerInstance = style.MustRegisterTimelineType(
		"Parser",
		"Logs of the DAG Processor Manager instance. Same DAG file can be parsed from multiple DAG Processor Manager instances at the same time thus this is shown as separated timelines.",
		"terminal",
		0.6,
		style.ColorWhite,
		style.ColorBlack,
		style.MustForceConvertSRGBHex("#F5F5F5"),
		style.ColorBlack,
		true,
		90,
		style.AlphabeticalSortPolicy(),
	)

	// TimelineTypeAirflowComponent is the style for Airflow components.
	TimelineTypeAirflowComponent = style.MustRegisterTimelineType(
		"Component",
		"Logs of the generic Airflow component",
		"extension",
		0.6,
		style.ColorWhite,
		style.ColorBlack,
		style.MustForceConvertSRGBHex("#F5F5F5"),
		style.ColorBlack,
		true,
		100,
		style.AlphabeticalSortPolicy(),
	)
)

The following block defines the registered timeline style TimelineTypes. These are registered as package-level variables so they are initialized immediately when this package is imported.

View Source
var (
	// TODO: This will be removed when the history pane starts showing states as well as verbs.
	VerbComposerTaskInstanceScheduled = style.MustRegisterVerb("Scheduled", style.MustForceConvertSRGBHex("#1E88E5"), style.ColorWhite, true)
	// TODO: This will be removed when the history pane starts showing states as well as verbs.
	VerbComposerTaskInstanceQueued = style.MustRegisterVerb("Queued", style.MustForceConvertSRGBHex("#22CC22"), style.ColorWhite, true)
	// TODO: This will be removed when the history pane starts showing states as well as verbs.
	VerbComposerTaskInstanceRunning = style.MustRegisterVerb("Running", style.MustForceConvertSRGBHex("#22CC22"), style.ColorWhite, true)
	// TODO: This will be removed when the history pane starts showing states as well as verbs.
	VerbComposerTaskInstanceUpForRetry = style.MustRegisterVerb("UpForRetry", style.MustForceConvertSRGBHex("#FF7700"), style.ColorWhite, true)
	// TODO: This will be removed when the history pane starts showing states as well as verbs.
	VerbComposerTaskInstanceSuccess = style.MustRegisterVerb("Success", style.MustForceConvertSRGBHex("#22CC22"), style.ColorWhite, true)
	// TODO: This will be removed when the history pane starts showing states as well as verbs.
	VerbComposerTaskInstanceFailed = style.MustRegisterVerb("Failed", style.MustForceConvertSRGBHex("#A51915"), style.ColorWhite, true)
	// TODO: This will be removed when the history pane starts showing states as well as verbs.
	VerbComposerTaskInstanceDeferred = style.MustRegisterVerb("Deferred", style.MustForceConvertSRGBHex("#9470DC"), style.ColorWhite, true)
	// TODO: This will be removed when the history pane starts showing states as well as verbs.
	VerbComposerTaskInstanceUpForReschedule = style.MustRegisterVerb("UpForReschedule", style.MustForceConvertSRGBHex("#FF7700"), style.ColorWhite, true)
	// TODO: This will be removed when the history pane starts showing states as well as verbs.
	VerbComposerTaskInstanceRemoved = style.MustRegisterVerb("Removed", style.MustForceConvertSRGBHex("#A51915"), style.ColorWhite, true)
	// TODO: This will be removed when the history pane starts showing states as well as verbs.
	VerbComposerTaskInstanceUpstreamFailed = style.MustRegisterVerb("UpstreamFailed", style.MustForceConvertSRGBHex("#A51915"), style.ColorWhite, true)
	// TODO: This will be removed when the history pane starts showing states as well as verbs.
	VerbComposerTaskInstanceZombie = style.MustRegisterVerb("Zombie", style.MustForceConvertSRGBHex("#696969"), style.ColorWhite, true)
	// TODO: This will be removed when the history pane starts showing states as well as verbs.
	VerbComposerTaskInstanceStats = style.MustRegisterVerb("Stats", style.MustForceConvertSRGBHex("#DDDDDD"), style.ColorWhite, true)
	// TODO: This will be removed when the history pane starts showing states as well as verbs.
	VerbComposerTaskInstanceUnimplemented = style.MustRegisterVerb("Unimplemented", style.MustForceConvertSRGBHex("#DDDDDD"), style.ColorWhite, true)
	// TODO: This will be removed when the history pane starts showing states as well as verbs.
	VerbComposerTaskInstanceSkipped = style.MustRegisterVerb("Skipped", style.MustForceConvertSRGBHex("#e60076"), style.ColorWhite, true)
)

The following block defines the registered timeline style Verbs. These are registered as package-level variables so they are initialized immediately when this package is imported.

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[struct{}](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[inspectiontaskbase.TimelineMapperResult](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[struct{}](GoogleCloudComposerTaskIDPrefix + "ingester-other")

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

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[struct{}](GoogleCloudComposerTaskIDPrefix + "ingester-scheduler")

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

View Source
var AirflowSchedulerLogToTimelineMapperTaskID = taskid.NewDefaultImplementationID[inspectiontaskbase.TimelineMapperResult](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[struct{}](GoogleCloudComposerTaskIDPrefix + "ingester-worker")

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

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:        "Managed Airflow",
	Description: "Gather and parse Managed Airflow logs to visualize environment operations on timelines.",
	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.

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.

View Source
var (
	LogTypeManagedAirflowEnvironment = style.MustRegisterLogType("managed airflow", "Managed Airflow Environment Logs", style.MustForceConvertSRGBHex("#88AA55"), style.ColorWhite)
)

The following block defines the registered timeline style LogTypes. These are registered as package-level variables so they are initialized immediately when this package is imported.

Functions

func ExtractComposerComponent added in v0.58.2

func ExtractComposerComponent(reader *structured.NodeReader) (string, error)

ExtractComposerComponent extracts the component name from a Composer log NodeReader.

func MustAirflowComponentTimeline added in v0.56.0

func MustAirflowComponentTimeline(ctx context.Context, envPath *khifilev6.TimelinePath, name string) *khifilev6.TimelinePath

MustAirflowComponentTimeline returns the timeline path for a specific Airflow component.

func MustAirflowComponentsRootTimeline added in v0.56.0

func MustAirflowComponentsRootTimeline(ctx context.Context, envPath *khifilev6.TimelinePath) *khifilev6.TimelinePath

MustAirflowComponentsRootTimeline returns the category root timeline path for Airflow components.

func MustAirflowDAGFileTimeline added in v0.56.0

func MustAirflowDAGFileTimeline(ctx context.Context, envPath *khifilev6.TimelinePath, filePath string) *khifilev6.TimelinePath

MustAirflowDAGFileTimeline returns the timeline path for a parsed DAG file (with GCS prefix trimmed).

func MustAirflowDAGFilesTimeline added in v0.57.4

func MustAirflowDAGFilesTimeline(ctx context.Context, envPath *khifilev6.TimelinePath) *khifilev6.TimelinePath

MustAirflowDAGFilesTimeline returns the root timeline path for DAG files.

func MustAirflowDAGProcessorManagerInstanceTimeline added in v0.56.0

func MustAirflowDAGProcessorManagerInstanceTimeline(ctx context.Context, envPath *khifilev6.TimelinePath, filePath, instanceID string) *khifilev6.TimelinePath

MustAirflowDAGProcessorManagerInstanceTimeline returns the timeline path for the manager instance that processed the file.

func MustAirflowDAGRunTimeline added in v0.56.0

func MustAirflowDAGRunTimeline(ctx context.Context, envPath *khifilev6.TimelinePath, dagID, runID string) *khifilev6.TimelinePath

MustAirflowDAGRunTimeline returns the timeline path for a specific Airflow DAG run.

func MustAirflowDAGTimeline added in v0.56.0

func MustAirflowDAGTimeline(ctx context.Context, envPath *khifilev6.TimelinePath, dagID string) *khifilev6.TimelinePath

MustAirflowDAGTimeline returns the timeline path for a specific Airflow DAG.

func MustAirflowDAGsRootTimeline added in v0.56.0

func MustAirflowDAGsRootTimeline(ctx context.Context, envPath *khifilev6.TimelinePath) *khifilev6.TimelinePath

MustAirflowDAGsRootTimeline returns the root timeline path for DAGs.

func MustAirflowTaskInstanceTimeline added in v0.56.0

func MustAirflowTaskInstanceTimeline(ctx context.Context, runPath *khifilev6.TimelinePath, taskName string) *khifilev6.TimelinePath

MustAirflowTaskInstanceTimeline returns the timeline path for a specific Airflow TaskInstance under the given runPath.

func MustAirflowTimeline added in v0.57.4

func MustAirflowTimeline(ctx context.Context, environmentName string) *khifilev6.TimelinePath

MustAirflowTimeline returns the root timeline path for an Airflow environment.

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) 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) ToYaml added in v0.52.8

func (a *AirflowWorker) ToYaml() string

type ComposerEnvironmentClusterFinder

type ComposerEnvironmentClusterFinder interface {
	GetGKEClusterNames(ctx context.Context, projectID, location, environment string, startTime, endTime time.Time) ([]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
}

ComposerFieldSet contains component identity info for Composer logs.

func ExtractComposer added in v0.58.2

func ExtractComposer(reader *structured.NodeReader) (ComposerFieldSet, error)

ExtractComposer extracts ComposerFieldSet from a NodeReader.

type ComposerTaskInstanceFieldSet added in v0.52.8

type ComposerTaskInstanceFieldSet struct {
	TaskInstance *AirflowTaskInstance
}

ComposerTaskInstanceFieldSet contains parsed Airflow task instance info.

func ExtractComposerTaskInstance added in v0.58.2

func ExtractComposerTaskInstance(reader *structured.NodeReader) (ComposerTaskInstanceFieldSet, error)

ExtractComposerTaskInstance extracts ComposerTaskInstanceFieldSet from a NodeReader.

type ComposerWorkerTaskInstanceFieldSet added in v0.52.8

type ComposerWorkerTaskInstanceFieldSet struct {
	TaskInstance *AirflowTaskInstance
}

ComposerWorkerTaskInstanceFieldSet contains worker-specific task instance info.

func ExtractComposerWorkerTaskInstance added in v0.58.2

func ExtractComposerWorkerTaskInstance(reader *structured.NodeReader) (ComposerWorkerTaskInstanceFieldSet, error)

ExtractComposerWorkerTaskInstance extracts ComposerWorkerTaskInstanceFieldSet from a NodeReader.

type EnvironmentClusterFinderImpl

type EnvironmentClusterFinderImpl struct{}

func (*EnvironmentClusterFinderImpl) GetGKEClusterNames added in v0.57.8

func (e *EnvironmentClusterFinderImpl) GetGKEClusterNames(ctx context.Context, projectID, location, environment string, startTime, endTime time.Time) ([]string, error)

GetGKEClusterNames implements ComposerEnvironmentClusterFinder.

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