Documentation
¶
Index ¶
- Constants
- Variables
- func MustAirflowComponentTimeline(ctx context.Context, envPath *khifilev6.TimelinePath, name string) *khifilev6.TimelinePath
- func MustAirflowComponentsRootTimeline(ctx context.Context, envPath *khifilev6.TimelinePath) *khifilev6.TimelinePath
- func MustAirflowDAGFileTimeline(ctx context.Context, envPath *khifilev6.TimelinePath, filePath string) *khifilev6.TimelinePath
- func MustAirflowDAGFilesTimeline(ctx context.Context, envPath *khifilev6.TimelinePath) *khifilev6.TimelinePath
- func MustAirflowDAGProcessorManagerInstanceTimeline(ctx context.Context, envPath *khifilev6.TimelinePath, ...) *khifilev6.TimelinePath
- func MustAirflowDAGRunTimeline(ctx context.Context, envPath *khifilev6.TimelinePath, dagID, runID string) *khifilev6.TimelinePath
- func MustAirflowDAGTimeline(ctx context.Context, envPath *khifilev6.TimelinePath, dagID string) *khifilev6.TimelinePath
- func MustAirflowDAGsRootTimeline(ctx context.Context, envPath *khifilev6.TimelinePath) *khifilev6.TimelinePath
- func MustAirflowTaskInstanceTimeline(ctx context.Context, runPath *khifilev6.TimelinePath, taskName string) *khifilev6.TimelinePath
- func MustAirflowTimeline(ctx context.Context, environmentName string) *khifilev6.TimelinePath
- type AirflowTaskInstance
- func (a *AirflowTaskInstance) DagId() string
- func (a *AirflowTaskInstance) Host() string
- func (a *AirflowTaskInstance) MapIndex() string
- func (a *AirflowTaskInstance) RunId() string
- func (a *AirflowTaskInstance) Status() Tistate
- func (a *AirflowTaskInstance) TaskId() string
- func (a *AirflowTaskInstance) ToYaml() string
- type AirflowWorker
- type ComposerEnvironmentClusterFinder
- type ComposerEnvironmentIdentity
- type ComposerEnvironmentListFetcher
- type ComposerEnvironmentListFetcherImpl
- type ComposerFieldSet
- type ComposerFieldSetReader
- type ComposerTaskInstanceFieldSet
- type ComposerTaskInstanceFieldSetReader
- type ComposerWorkerTaskInstanceFieldSet
- type ComposerWorkerTaskInstanceFieldSetReader
- type EnvironmentClusterFinderImpl
- type Tistate
Constants ¶
const InspectionTypeID = "gcp-composer"
InspectionTypeID is the inspection type id for google cloud composer.
Variables ¶
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.
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.
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.
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.
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.
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.
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.
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.
var AirflowOtherLogFilterTaskID taskid.TaskImplementationID[[]*log.Log] = taskid.NewDefaultImplementationID[[]*log.Log](GoogleCloudComposerTaskIDPrefix + "filter-other")
AirflowOtherLogFilterTaskID is the task id for filtering other Airflow logs.
var AirflowOtherLogGrouperTaskID = taskid.NewDefaultImplementationID[inspectiontaskbase.LogGroupMap](GoogleCloudComposerTaskIDPrefix + "grouper-other")
AirflowOtherLogGrouperTaskID is the task id for the task that groups other Airflow logs.
var AirflowOtherLogIngesterTaskID = taskid.NewDefaultImplementationID[[]*log.Log](GoogleCloudComposerTaskIDPrefix + "ingester-other")
AirflowOtherLogIngesterTaskID is the task id for the task that ingests other Airflow logs.
var AirflowOtherLogToTimelineMapperTaskID = taskid.NewDefaultImplementationID[inspectiontaskbase.TimelineMapperResult](GoogleCloudComposerTaskIDPrefix + "mapper-other")
AirflowOtherLogToTimelineMapperTaskID is the task id for the task that maps other Airflow logs to timeline events.
var AirflowSchedulerLogFilterTaskID taskid.TaskImplementationID[[]*log.Log] = taskid.NewDefaultImplementationID[[]*log.Log](GoogleCloudComposerTaskIDPrefix + "filter-scheduler")
AirflowSchedulerLogFilterTaskID is the task id for filtering Airflow scheduler logs.
var AirflowSchedulerLogGrouperTaskID = taskid.NewDefaultImplementationID[inspectiontaskbase.LogGroupMap](GoogleCloudComposerTaskIDPrefix + "grouper-scheduler")
AirflowSchedulerLogGrouperTaskID is the task id for the task that groups Airflow scheduler logs.
var AirflowSchedulerLogIngesterTaskID = taskid.NewDefaultImplementationID[[]*log.Log](GoogleCloudComposerTaskIDPrefix + "ingester-scheduler")
AirflowSchedulerLogIngesterTaskID is the task id for the task that ingests Airflow scheduler logs.
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.
var AirflowWorkerLogFilterTaskID taskid.TaskImplementationID[[]*log.Log] = taskid.NewDefaultImplementationID[[]*log.Log](GoogleCloudComposerTaskIDPrefix + "filter-worker")
AirflowWorkerLogFilterTaskID is the task id for filtering Airflow worker logs.
var AirflowWorkerLogGrouperTaskID = taskid.NewDefaultImplementationID[inspectiontaskbase.LogGroupMap](GoogleCloudComposerTaskIDPrefix + "grouper-worker")
AirflowWorkerLogGrouperTaskID is the task id for the task that groups Airflow worker logs.
var AirflowWorkerLogIngesterTaskID = taskid.NewDefaultImplementationID[[]*log.Log](GoogleCloudComposerTaskIDPrefix + "ingester-worker")
AirflowWorkerLogIngesterTaskID is the task id for the task that ingests Airflow worker logs.
var AirflowWorkerLogToTimelineMapperTaskID = taskid.NewDefaultImplementationID[inspectiontaskbase.TimelineMapperResult](GoogleCloudComposerTaskIDPrefix + "mapper-worker")
AirflowWorkerLogToTimelineMapperTaskID is the task id for the task that maps Airflow worker logs to timeline events.
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.
var AutocompleteComposerComponentsTaskID = taskid.NewDefaultImplementationID[*inspectioncore_contract.AutocompleteResult[string]](GoogleCloudComposerTaskIDPrefix + "autocomplete/composer-components")
AutocompleteComposerComponentsTaskID is the task id for autocompleting component names from Cloud Monitoring.
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.
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.
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.
var ClusterIdentityTaskID = taskid.NewDefaultImplementationID[googlecloudk8scommon_contract.GoogleCloudClusterIdentity](GoogleCloudComposerTaskIDPrefix + "cluster-identity")
ClusterIdentityTaskID is the task id for aliasing the cluster identity.
var ComposerClusterNamePrefixTaskID = taskid.NewImplementationID(googlecloudk8scommon_contract.ClusterNamePrefixTaskRef, "composer")
ComposerClusterNamePrefixTaskID is the task id for the task that returns the GKE cluster name prefix used by Cloud Composer.
var ComposerEnvironmentClusterFinderTaskID = taskid.NewDefaultImplementationID[ComposerEnvironmentClusterFinder](GoogleCloudComposerTaskIDPrefix + "composer-environment-cluster-finder")
ComposerEnvironmentClusterFinderTaskID is the task id for injecting ComposerEnvironmentClusterFinder instance.
var ComposerEnvironmentListFetcherTaskID = taskid.NewDefaultImplementationID[ComposerEnvironmentListFetcher](GoogleCloudComposerTaskIDPrefix + "composer-environment-list-fetcher")
ComposerEnvironmentListFetcherTaskID is the task id for injecting ComposerEnvironmentListFetcher instance.
var ComposerInspectionType = coreinspection.InspectionType{ Id: InspectionTypeID, Name: "Managed Service for Apache Airflow", Description: "Gather and parse Managed Service for Apache Airflow environment logs (supporting GKE logs in Managed Airflow v2 and Apache Airflow logs (version 2.0.0 or higher) in any Managed Airflow version) 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.
var ComposerLogsFieldSetReadTaskID taskid.TaskImplementationID[[]*log.Log] = taskid.NewDefaultImplementationID[[]*log.Log](GoogleCloudComposerTaskIDPrefix + "fieldsetread")
ComposerLogsFieldSetReadTaskID is the task id for the task that reads fieldsets from composer logs.
var ComposerLogsQueryTaskID taskid.TaskImplementationID[[]*log.Log] = taskid.NewDefaultImplementationID[[]*log.Log](GoogleCloudComposerTaskIDPrefix + "query-composer-logs")
ComposerLogsQueryTaskID is the task id for the task that queries Logs from Cloud Logging.
var ComposerLogsTailTaskID = taskid.NewDefaultImplementationID[struct{}](GoogleCloudComposerTaskIDPrefix + "tail-composer-logs")
ComposerLogsTailTaskID is the task id for unifying composer logs feature.
var ErrEnvironmentClusterNotFound = errors.New("not found")
var GoogleCloudComposerTaskIDPrefix = "cloud.google.com/composer/"
GoogleCloudComposerTaskIDPrefix is the prefix for all task ids related to google cloud composer.
var InputComposerComponentsTaskID taskid.TaskImplementationID[[]string] = taskid.NewDefaultImplementationID[[]string](GoogleCloudComposerTaskIDPrefix + "input/composer/components")
InputComposerComponentsTaskID is the task id for selecting target Composer components.
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.
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 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 (*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 ComposerEnvironmentIdentity ¶ added in v0.52.0
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
func (c *ComposerFieldSetReader) Read(reader *structured.NodeReader) (log.FieldSet, error)
type ComposerTaskInstanceFieldSet ¶ added in v0.52.8
type ComposerTaskInstanceFieldSet struct {
TaskInstance *AirflowTaskInstance
}
func (*ComposerTaskInstanceFieldSet) Kind ¶ added in v0.52.8
func (c *ComposerTaskInstanceFieldSet) Kind() string
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
func (c *ComposerTaskInstanceFieldSetReader) Read(reader *structured.NodeReader) (log.FieldSet, error)
type ComposerWorkerTaskInstanceFieldSet ¶ added in v0.52.8
type ComposerWorkerTaskInstanceFieldSet struct {
TaskInstance *AirflowTaskInstance
}
func (*ComposerWorkerTaskInstanceFieldSet) Kind ¶ added in v0.52.8
func (c *ComposerWorkerTaskInstanceFieldSet) Kind() string
type ComposerWorkerTaskInstanceFieldSetReader ¶ added in v0.52.8
type ComposerWorkerTaskInstanceFieldSetReader struct{}
func (*ComposerWorkerTaskInstanceFieldSetReader) FieldSetKind ¶ added in v0.52.8
func (c *ComposerWorkerTaskInstanceFieldSetReader) FieldSetKind() string
func (*ComposerWorkerTaskInstanceFieldSetReader) Read ¶ added in v0.52.8
func (c *ComposerWorkerTaskInstanceFieldSetReader) Read(reader *structured.NodeReader) (log.FieldSet, error)
Read parses structured log data to extract worker TaskInstance information.
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" )