Documentation
¶
Index ¶
- Constants
- Variables
- func ExtractGCPMainMessage(reader *structured.NodeReader) (string, error)
- func ExtractGCPSeverity(reader *structured.NodeReader) (*pb.Severity, error)
- func MustGCPOperationTimeline(ctx context.Context, parentTimeline *khifilev6.TimelinePath, ...) *khifilev6.TimelinePath
- func MustGCPProjectTimeline(ctx context.Context, projectID string) *khifilev6.TimelinePath
- func MustGCPResourceTimeline(ctx context.Context, resourceTypePath *khifilev6.TimelinePath, ...) *khifilev6.TimelinePath
- func MustGCPResourceTypeTimeline(ctx context.Context, projectPath *khifilev6.TimelinePath, resourceType string) *khifilev6.TimelinePath
- func MustGKEClusterTimeline(ctx context.Context, projectPath *khifilev6.TimelinePath, clusterName string) *khifilev6.TimelinePath
- func MustGKENodePoolTimeline(ctx context.Context, gkeClusterTimeline *khifilev6.TimelinePath, ...) *khifilev6.TimelinePath
- func MustManagedAirflowEnvironmentTimeline(ctx context.Context, projectPath *khifilev6.TimelinePath, ...) *khifilev6.TimelinePath
- func NewGCPOperationLogIngester(rawLogTask taskid.TaskReference[[]*log.Log], logType *pb.LogType) inspectiontaskbase.LogIngester
- func NewGCPOperationLogIngesterTask(taskID taskid.TaskImplementationID[struct{}], ...) coretask.Task[struct{}]
- func NewListLogEntriesTask(taskSetting ListLogEntriesTaskSetting) coretask.Task[[]*log.Log]
- func NewStructuredListLogEntriesTask(taskSetting StructuredListLogEntriesTaskSetting) coretask.Task[[]*log.Log]
- func ParseGCPSeverity(gcpSeverity string) *pb.Severity
- func ProcessGCPClusterNodepoolOperationLog(ctx context.Context, cs *khifilev6.TimelineChangeSet, ...)
- type GCPAccessLogFieldSet
- type GCPAuditLogFieldSet
- func (g *GCPAuditLogFieldSet) Ending() bool
- func (g *GCPAuditLogFieldSet) GuessRevisionVerb() *pb.Verb
- func (g *GCPAuditLogFieldSet) ImmediateOperation() bool
- func (g *GCPAuditLogFieldSet) RequestString() (string, error)
- func (g *GCPAuditLogFieldSet) ResponseString() (string, error)
- func (g *GCPAuditLogFieldSet) Starting() bool
- type GCPMainMessageFieldSet
- type GCPOperationLogIngester
- type GCPOperationTracker
- func (t *GCPOperationTracker) HasResourceRevision(path *khifilev6.TimelinePath) bool
- func (t *GCPOperationTracker) HasStarted(operationID string) bool
- func (t *GCPOperationTracker) MarkResourceRevision(path *khifilev6.TimelinePath)
- func (t *GCPOperationTracker) MarkStarted(operationID string)
- func (t *GCPOperationTracker) ProcessOperationLog(ctx context.Context, cs *khifilev6.TimelineChangeSet, ...)
- func (t *GCPOperationTracker) TrackAndGetManifest(audit *GCPAuditLogFieldSet) (structured.Node, bool)
- type ListLogEntriesTaskDescription
- type ListLogEntriesTaskSetting
- type LocationFetcher
- type LogFetchProgress
- type LogFetcher
- type ProgressReportableLogFetcher
- type QueryResourceNames
- type ResourceNamesInput
- type StandardProgressReportableLogFetcher
- type StructuredListLogEntriesTaskSetting
- type TimePartitioningProgressReportableLogFetcher
Constants ¶
const ( // FormBasePriority is the base priority for Google Cloud common forms. FormBasePriority = 100000 // PriorityForQueryTimeGroup is the priority for the query time group. PriorityForQueryTimeGroup = FormBasePriority + 50000 // PriorityForResourceIdentifierGroup is the priority for the resource identifier group. PriorityForResourceIdentifierGroup = FormBasePriority + 40000 // PriorityForK8sResourceFilterGroup is the priority for the k8s resource filter group. PriorityForK8sResourceFilterGroup = FormBasePriority + 30000 )
const ( // InspectionTypeLabelKeyProduct is the label key for the Google Cloud product name running on the cluster. // Expected values of this label key: "composer", etc. InspectionTypeLabelKeyProduct = "cloud.google.com/product" // InspectionTypeLabelKeyClusterType is the label key for the cluster type. // Expected values of this label key: "gke", "gdc", "gke_multicloud", etc. InspectionTypeLabelKeyClusterType = "cloud.google.com/cluster_type" // InspectionTypeLabelKeyClusterSubType is the label key for the cluster sub type. // Expected values of this label key: "aws", "azure", etc. InspectionTypeLabelKeyClusterSubType = "cloud.google.com/cluster_subtype" )
Variables ¶
var ( RevisionStateOperationStarted = style.MustRegisterRevisionState( "Processing operation", "change_circle", "The GCP API long-running operation is currently in progress.", style.Color{R: 0.012, G: 0.671, B: 0.012, A: 1.0}, pb.RevisionStateStyle_REVISION_STATE_STYLE_NORMAL, ) RevisionStateOperationSucceed = style.MustRegisterRevisionState( "Operation succeeded", "check_circle", "The GCP API long-running operation has completed successfully.", style.Color{R: 0.812, G: 0.812, B: 0.812, A: 1.0}, pb.RevisionStateStyle_REVISION_STATE_STYLE_DELETED, ) RevisionStateOperationFailed = style.MustRegisterRevisionState( "Operation failed", "error", "The GCP API long-running operation has failed.", style.Color{R: 1.000, G: 0.000, B: 0.000, A: 1.0}, pb.RevisionStateStyle_REVISION_STATE_STYLE_DELETED, ) RevisionStateOperationStartedLogNotFound = style.MustRegisterRevisionState( "Operation started, but starting log not found", "unknown_document", "The operation was started, but the starting log entry was not found in the selected time range. Try adjusting the time range.", style.Color{R: 0.012, G: 0.671, B: 0.012, 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 ( // TimelineTypeManagedAirflowEnvironment is the timeline style for a Managed Airflow environment under a GCP project. TimelineTypeManagedAirflowEnvironment = style.MustRegisterTimelineType( "Managed Airflow", "Timeline for Managed Airflow", "cloud", 1.0, style.Color{R: 0.780, G: 0.863, B: 1.000, A: 1.0}, style.ColorBlack, style.Color{R: 0.780, G: 0.863, B: 1.000, A: 1.0}, style.ColorBlack, true, 40, style.AlphabeticalSortPolicy(), ) TimelineTypeGKE = style.MustRegisterTimelineType( "gke", "Control plane operations and lifecycle logs of the GKE cluster", "cloud", 1.0, style.Color{R: 0.780, G: 0.863, B: 1.000, A: 1.0}, style.ColorBlack, style.Color{R: 0.780, G: 0.863, B: 1.000, A: 1.0}, style.ColorBlack, true, 40, style.AlphabeticalSortPolicy(), ) // TimelineTypeGKEControlPlanes is the timeline style for a GKE control planes folder. TimelineTypeGKEControlPlanes = style.MustRegisterTimelineType( "controlplanes", "Control Plane", "category", 0.6, style.Color{R: 0.361, G: 0.604, B: 1.000, A: 1.0}, style.ColorWhite, style.Color{R: 0.361, G: 0.604, B: 1.000, A: 1.0}, style.ColorWhite, false, 3002, style.AlphabeticalSortPolicy(), ) // TimelineTypeGKENodePools is the timeline style for a GKE node pools folder. TimelineTypeGKENodePools = style.MustRegisterTimelineType( "nodepools", "Node Pools", "dns", 0.6, style.Color{R: 0.361, G: 0.604, B: 1.000, A: 1.0}, style.ColorWhite, style.Color{R: 0.361, G: 0.604, B: 1.000, A: 1.0}, style.ColorWhite, false, 3001, style.AlphabeticalSortPolicy(), ) // TimelineTypeOtherGKEResources is the timeline style for other GKE resources. TimelineTypeOtherGKEResources = style.MustRegisterTimelineType( "other_gke_resources", "Other GKE Resources", "category", 0.6, style.Color{R: 0.361, G: 0.604, B: 1.000, A: 1.0}, style.ColorWhite, style.Color{R: 0.361, G: 0.604, B: 1.000, A: 1.0}, style.ColorWhite, false, 3003, style.AlphabeticalSortPolicy(), ) // TimelineTypeGKENodePool is the style for a GKE node pool. TimelineTypeGKENodePool = style.MustRegisterTimelineType( "nodepool", "Grouping timeline for GKE nodepools", "dns", 0.8, style.Color{R: 0.941, G: 0.965, B: 1.000, A: 1.0}, style.ColorBlack, style.Color{R: 0.941, G: 0.965, B: 1.000, A: 1.0}, style.ColorBlack, true, 45, style.AlphabeticalSortPolicy(), ) TimelineTypeOperation = style.MustRegisterTimelineType( "operation", "Google Cloud operations associated with the resource", "engineering", 0.6, style.ColorWhite, style.ColorBlack, style.ColorBlack, style.ColorWhite, true, 3000, style.ChronologicalSortPolicy(0), ) TimelineTypeGCPProject = style.MustRegisterTimelineType( "project", "Timeline representing a Google Cloud project", "cloud", 1.0, style.Color{R: 0.102, G: 0.451, B: 0.910, A: 1.0}, style.ColorWhite, style.Color{R: 0.102, G: 0.451, B: 0.910, A: 1.0}, style.ColorWhite, true, 30, style.AlphabeticalSortPolicy(), ) TimelineTypeGCPResourceType = style.MustRegisterTimelineType( "gcp_resource_type", "Grouping timeline for Google Cloud resource types", "category", 0.6, style.ColorWhite, style.ColorBlack, style.MustForceConvertSRGBHex("#34A853"), style.ColorWhite, true, 31, style.AlphabeticalSortPolicy(), ) TimelineTypeGCPResource = style.MustRegisterTimelineType( "gcp_resource", "Timeline representing a Google Cloud resource", "deployed_code", 0.6, style.ColorWhite, style.ColorBlack, style.MustForceConvertSRGBHex("#FBBC05"), style.ColorBlack, true, 32, 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 ( VerbOperationStart = style.MustRegisterVerb("Start", style.MustForceConvertSRGBHex("#22CC22"), style.ColorWhite, true) VerbOperationFinish = style.MustRegisterVerb("Finish", style.MustForceConvertSRGBHex("#9999CC"), 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 APICallOptionsInjectorContextKey = typedmap.NewTypedKey[*[]googlecloud.CallOptionInjectorOption]("api-call-option-injector-options")
APICallOptionsInjectorContextKey is the key to retrieve the list of googlecloud.CallOptionInjectorOption from task context. The value is injected on the task server during the initialization.
var APIClientCallOptionsInjectorTaskID = taskid.NewDefaultImplementationID[*googlecloud.CallOptionInjector](GoogleCloudCommonTaskIDPrefix + "api-client-option-injector")
APIClientCallOptionsInjectorTaskID is the task ID to inject CallOptionInjector reference.
var APIClientFactoryOptionsContextKey = typedmap.NewTypedKey[*[]googlecloud.ClientFactoryOption]("api-client-factory-options")
APIClientFactoryOptionsContextKey is the key to retrieve googlecloud.ClientFactoryOption from task context. The value is injected on the task server during the initialization.
var APIClientFactoryOptionsTaskID = taskid.NewDefaultImplementationID[[]googlecloud.ClientFactoryOption](GoogleCloudCommonTaskIDPrefix + "api-client-factory-options")
APIClientFactoryOptionsTaskID is the task ID to generate options list for the ClientFactory. This can be overridden with the selection priority label defined in the coretask package.
var APIClientFactoryTaskID = taskid.NewDefaultImplementationID[*googlecloud.ClientFactory](GoogleCloudCommonTaskIDPrefix + "api-client-factory")
APIClientFactoryTaskID is the task ID to generate the ClientFactory. This factory is instantiated with the options generated from the task with APIClientFactoryOptionsTaskID.
var AutocompleteLocationTaskID taskid.TaskImplementationID[*inspectioncore_contract.AutocompleteResult[string]] = taskid.NewDefaultImplementationID[*inspectioncore_contract.AutocompleteResult[string]](GoogleCloudCommonTaskIDPrefix + "autocomplete-location")
AutocompleteLocationTaskID is the task ID for the location autocomplete.
var DefaultAPIClientOptionTasksPriority = 1000
DefaultAPIClientOptionTasksPriority is the selection priority of the default implementation for APIClientFactoryOptionsTaskID and APICallOptionsInjectorTask. Users can define another task for the task ID with higher priority to override the options.
var GCPAuditLogCacheKey = structured.NewCacheKey[*GCPAuditLogFieldSet]()
GCPAuditLogCacheKey identifies the cached GCPAuditLogFieldSet on a NodeReader.
var GCPSeverityCacheKey = structured.NewCacheKey[*pb.Severity]()
GCPSeverityCacheKey identifies the cached severity on a NodeReader.
var GoogleCloudCommonTaskIDPrefix = "cloud.google.com/common/"
GoogleCloudCommonTaskIDPrefix is the prefix for Google Cloud common task IDs.
var InputDurationTaskID = taskid.NewDefaultImplementationID[time.Duration](GoogleCloudCommonTaskIDPrefix + "input-duration")
InputDurationTaskID is the task ID for the duration of the log query.
var InputEndTimeTaskID = taskid.NewDefaultImplementationID[time.Time](GoogleCloudCommonTaskIDPrefix + "input-end-time")
InputEndTimeTaskID is the task ID for the end time of the log query.
var InputLocationsTaskID = taskid.NewDefaultImplementationID[string](GoogleCloudCommonTaskIDPrefix + "input-location")
InputLocationsTaskID is the task ID for the locations of the target resource.
var InputLoggingFilterResourceNameTaskID = taskid.NewDefaultImplementationID[*ResourceNamesInput](GoogleCloudCommonTaskIDPrefix + "input-logging-filter-resource-name")
InputLoggingFilterResourceNameTaskID is the task ID to get log query target resource names.
var InputProjectIdTaskID = taskid.NewDefaultImplementationID[string](GoogleCloudCommonTaskIDPrefix + "input-project-id")
InputProjectIdTaskID is the task ID for the Google Cloud project ID.
var InputStartTimeTaskID = taskid.NewDefaultImplementationID[time.Time](GoogleCloudCommonTaskIDPrefix + "input-start-time")
InputStartTimeTaskID is the task ID for the start time of the log query. This is computed from InputDurationTask and InputEndTimeTask.
var LocationFetcherTaskID = taskid.NewDefaultImplementationID[LocationFetcher](GoogleCloudCommonTaskIDPrefix + "location-fetcher")
LocationFetcherTaskID is the task ID to inject the instance of LocationFetcher.
var LogEstimatorCacheKey = typedmap.NewTypedKey[*logestimator.CachedStructuredLogEstimator]("googlecloud.logestimator.cache")
LogEstimatorCacheKey is the key to retrieve or store CachedStructuredLogEstimator in the InspectionSharedMap.
var LoggingFetcherTaskID = taskid.NewDefaultImplementationID[LogFetcher](GoogleCloudCommonTaskIDPrefix + "log-fetcher")
LoggingFetcherTaskID is the task ID to inject the instance of LogFetcher.
var RequestOptionalInputResourceNameTaskLabel = typedmap.NewTypedKey[string]("request-optional-input-resource-name")
RequestOptionalInputResourceNameTaskLabel is a label assigned to a task that requests the Cloud Logging resource name optionally. The value is the query ID.
Functions ¶
func ExtractGCPMainMessage ¶ added in v0.58.2
func ExtractGCPMainMessage(reader *structured.NodeReader) (string, error)
ExtractGCPMainMessage reads main message from the content of log stored on Cloud Logging. It treats fields as its main message in order: protoPayload > textPayload > jsonPayload.**** > jsonPayload > labels.
func ExtractGCPSeverity ¶ added in v0.58.2
func ExtractGCPSeverity(reader *structured.NodeReader) (*pb.Severity, error)
ExtractGCPSeverity extracts severity from a GCP Cloud Logging entry.
func MustGCPOperationTimeline ¶ added in v0.56.0
func MustGCPOperationTimeline(ctx context.Context, parentTimeline *khifilev6.TimelinePath, shortMethodName string, operationID string) *khifilev6.TimelinePath
MustGCPOperationTimeline returns the timeline path for a GCP long running operation. The operation timeline is nested under its associated GKE or GCP resource timeline path.
func MustGCPProjectTimeline ¶ added in v0.56.0
func MustGCPProjectTimeline(ctx context.Context, projectID string) *khifilev6.TimelinePath
MustGCPProjectTimeline returns the timeline path for a Google Cloud Project root timeline.
func MustGCPResourceTimeline ¶ added in v0.56.0
func MustGCPResourceTimeline(ctx context.Context, resourceTypePath *khifilev6.TimelinePath, resourceName string) *khifilev6.TimelinePath
MustGCPResourceTimeline returns the timeline path for a GCP resource under a resource type.
func MustGCPResourceTypeTimeline ¶ added in v0.56.0
func MustGCPResourceTypeTimeline(ctx context.Context, projectPath *khifilev6.TimelinePath, resourceType string) *khifilev6.TimelinePath
MustGCPResourceTypeTimeline returns the timeline path for a GCP resource type layer under a GCP Project.
func MustGKEClusterTimeline ¶ added in v0.56.0
func MustGKEClusterTimeline(ctx context.Context, projectPath *khifilev6.TimelinePath, clusterName string) *khifilev6.TimelinePath
MustGKEClusterTimeline returns the timeline path for a GKE Cluster under a GCP Project.
func MustGKENodePoolTimeline ¶ added in v0.56.0
func MustGKENodePoolTimeline(ctx context.Context, gkeClusterTimeline *khifilev6.TimelinePath, nodePoolName string) *khifilev6.TimelinePath
MustGKENodePoolTimeline returns the timeline path for a GKE NodePool under a GKE Cluster.
func MustManagedAirflowEnvironmentTimeline ¶ added in v0.57.7
func MustManagedAirflowEnvironmentTimeline(ctx context.Context, projectPath *khifilev6.TimelinePath, environmentName string) *khifilev6.TimelinePath
MustManagedAirflowEnvironmentTimeline returns the timeline path for a Cloud Composer Environment under a GCP Project.
func NewGCPOperationLogIngester ¶ added in v0.56.1
func NewGCPOperationLogIngester(rawLogTask taskid.TaskReference[[]*log.Log], logType *pb.LogType) inspectiontaskbase.LogIngester
NewGCPOperationLogIngester creates a new GCPOperationLogIngester.
func NewGCPOperationLogIngesterTask ¶ added in v0.56.1
func NewGCPOperationLogIngesterTask(taskID taskid.TaskImplementationID[struct{}], rawLogTask taskid.TaskReference[[]*log.Log], logType *pb.LogType) coretask.Task[struct{}]
NewGCPOperationLogIngesterTask returns a new log ingester task for GCP Operation audit logs.
func NewListLogEntriesTask ¶
func NewListLogEntriesTask(taskSetting ListLogEntriesTaskSetting) coretask.Task[[]*log.Log]
NewListLogEntriesTask creates a new task that lists log entries from Cloud Logging based on the provided settings.
func NewStructuredListLogEntriesTask ¶ added in v0.58.2
func NewStructuredListLogEntriesTask(taskSetting StructuredListLogEntriesTaskSetting) coretask.Task[[]*log.Log]
NewStructuredListLogEntriesTask creates a new task that queries logs from Cloud Logging using StructuredLogQuery. In DryRun mode, it estimates log volumes and populates QueryMetadata with estimated counts.
func ParseGCPSeverity ¶ added in v0.56.0
ParseGCPSeverity converts a GCP Cloud Logging severity string into a timeline style Severity. It maps the GCP log severities defined in https://cloud.google.com/logging/docs/reference/v2/rest/v2/LogEntry#logseverity to KHI's registered timeline style Severities in the core contract.
func ProcessGCPClusterNodepoolOperationLog ¶ added in v0.56.2
func ProcessGCPClusterNodepoolOperationLog( ctx context.Context, cs *khifilev6.TimelineChangeSet, tracker *GCPOperationTracker, targetTimeline *khifilev6.TimelinePath, operationTimeline *khifilev6.TimelinePath, audit *GCPAuditLogFieldSet, logTimestamp time.Time, shortMethodName string, isCluster bool, )
// ProcessGCPClusterNodepoolOperationLog processes a GCP operation log for a cluster or node pool resource timeline. It generates resource creation/deletion/enrollment/unenrollment revisions, handles missing start logs by prepending appropriate LogNotFound revisions at Unix time 0, and updates operation tracking.
Types ¶
type GCPAccessLogFieldSet ¶
type GCPAccessLogFieldSet struct {
Method string
RequestURL string
RequestSize int64
Status int
ResponseSize int64
UserAgent string
RemoteIP string
ServerIP string
Referer string
Latency string
Protocol string
}
GCPAccessLogFieldSet represents HTTP access log fields from Cloud Logging.
func ExtractGCPAccessLog ¶ added in v0.58.2
func ExtractGCPAccessLog(reader *structured.NodeReader) (GCPAccessLogFieldSet, error)
ExtractGCPAccessLog extracts GCP Access Log fields from a NodeReader.
type GCPAuditLogFieldSet ¶
type GCPAuditLogFieldSet struct {
ProjectID string
OperationID string
OperationFirst bool
OperationLast bool
MethodName string
ResourceName string
PrincipalEmail string
Status int
StatusMessage string
Request *structured.NodeReader
Response *structured.NodeReader
}
GCPAuditLogFieldSet represents the parsed fields of a GCP Cloud Audit Log entry.
func ExtractGCPAuditLog ¶ added in v0.58.2
func ExtractGCPAuditLog(reader *structured.NodeReader) (GCPAuditLogFieldSet, error)
ExtractGCPAuditLog extracts GCP Audit Log fields from a NodeReader.
func (*GCPAuditLogFieldSet) Ending ¶
func (g *GCPAuditLogFieldSet) Ending() bool
Ending returns true when the operation is long running operation and the log entry is for the ending timing.
func (*GCPAuditLogFieldSet) GuessRevisionVerb ¶ added in v0.54.0
func (g *GCPAuditLogFieldSet) GuessRevisionVerb() *pb.Verb
GuessRevisionVerb returns the guessed revision verb from the method name.
func (*GCPAuditLogFieldSet) ImmediateOperation ¶
func (g *GCPAuditLogFieldSet) ImmediateOperation() bool
ImmediateOperation returns true when the log represents an operation completes immediately.
func (*GCPAuditLogFieldSet) RequestString ¶
func (g *GCPAuditLogFieldSet) RequestString() (string, error)
RequestString returns the request body as a YAML string.
func (*GCPAuditLogFieldSet) ResponseString ¶
func (g *GCPAuditLogFieldSet) ResponseString() (string, error)
ResponseString returns the response body as a YAML string.
func (*GCPAuditLogFieldSet) Starting ¶
func (g *GCPAuditLogFieldSet) Starting() bool
Starting returns true when the operation is long running operation and the log entry is for the starting timing.
type GCPMainMessageFieldSet ¶ added in v0.56.0
type GCPMainMessageFieldSet struct {
MainMessage string
}
GCPMainMessageFieldSet represents the main message parsed from a GCP log.
type GCPOperationLogIngester ¶ added in v0.56.1
type GCPOperationLogIngester struct {
// contains filtered or unexported fields
}
GCPOperationLogIngester is a common LogIngester implementation for GCP Operation audit logs.
func (*GCPOperationLogIngester) Dependencies ¶ added in v0.56.1
func (i *GCPOperationLogIngester) Dependencies() []taskid.UntypedTaskReference
Dependencies returns additional task dependencies of the ingester.
func (*GCPOperationLogIngester) ProcessLog ¶ added in v0.56.1
func (i *GCPOperationLogIngester) ProcessLog(ctx context.Context, l *log.Log) (*khifilev6.LogChangeSet, error)
ProcessLog parses raw log entry and populates the LogChangeSet.
func (*GCPOperationLogIngester) RawLogTask ¶ added in v0.56.1
func (i *GCPOperationLogIngester) RawLogTask() taskid.TaskReference[[]*log.Log]
RawLogTask returns the task reference that provides the raw logs to ingest.
type GCPOperationTracker ¶ added in v0.56.1
type GCPOperationTracker struct {
// contains filtered or unexported fields
}
GCPOperationTracker tracks operation start/finish logs within a group and generates revisions.
func NewGCPOperationTracker ¶ added in v0.56.1
func NewGCPOperationTracker() *GCPOperationTracker
NewGCPOperationTracker creates a new GCPOperationTracker.
func (*GCPOperationTracker) HasResourceRevision ¶ added in v0.56.2
func (t *GCPOperationTracker) HasResourceRevision(path *khifilev6.TimelinePath) bool
HasResourceRevision returns true if any revision has been added to the given resource timeline path.
func (*GCPOperationTracker) HasStarted ¶ added in v0.56.2
func (t *GCPOperationTracker) HasStarted(operationID string) bool
HasStarted returns true if the operation start log for the given operation ID was observed.
func (*GCPOperationTracker) MarkResourceRevision ¶ added in v0.56.2
func (t *GCPOperationTracker) MarkResourceRevision(path *khifilev6.TimelinePath)
MarkResourceRevision records that a revision was added to the given resource timeline path.
func (*GCPOperationTracker) MarkStarted ¶ added in v0.57.7
func (t *GCPOperationTracker) MarkStarted(operationID string)
MarkStarted records that the operation start log for the given operation ID was observed.
func (*GCPOperationTracker) ProcessOperationLog ¶ added in v0.56.1
func (t *GCPOperationTracker) ProcessOperationLog(ctx context.Context, cs *khifilev6.TimelineChangeSet, targetPath *khifilev6.TimelinePath, audit *GCPAuditLogFieldSet, timestamp time.Time)
ProcessOperationLog adds necessary operation revisions or events to the TimelineChangeSet. If an ending log is encountered without a prior starting log, it automatically prepends a dummy starting revision with RevisionStateOperationStartedLogNotFound at Unix time 0.
func (*GCPOperationTracker) TrackAndGetManifest ¶ added in v0.56.1
func (t *GCPOperationTracker) TrackAndGetManifest(audit *GCPAuditLogFieldSet) (structured.Node, bool)
TrackAndGetManifest tracks the latest resource manifest from an audit log and returns it if updated.
type ListLogEntriesTaskDescription ¶
ListLogEntriesTaskDescription holds descriptive information for a task to list log entries from CloudLogging.
type ListLogEntriesTaskSetting ¶
type ListLogEntriesTaskSetting interface {
// TaskID returns the task ID for the Cloud Logging list log entries task.
TaskID() taskid.TaskImplementationID[[]*log.Log]
// Dependencies returns the list of dependencies for the Cloud Logging list log entries task.
// Return the dependency task reference IDs when the result is used in DefaultResourceNames(), LogFilters() or TimePartitionCount().
Dependencies() []taskid.UntypedTaskReference
// DefaultResourceNames returns the list of resource names for the Cloud Logging list log entries task.
// This is just a default value for the resource name. Users can override this value with the form field.
// Return the list of resource names. ref: https://cloud.google.com/logging/docs/reference/v2/rest/v2/entries/list
DefaultResourceNames(ctx context.Context) ([]string, error)
// LogFilters returns the list of log filters for the Cloud Logging list log entries task.
// When generated logging filter can exceed the 20,000 character maximum limit in Cloud Logging, return multiple subset query.
// Result includes the logs for all log filters.
LogFilters(ctx context.Context, taskMode inspectioncore_contract.InspectionTaskModeType) ([]string, error)
// TimePartitionCount returns the number of time partitions for the Cloud Logging list log entries task.
// ListLogEntriesTask split the duration into the number of partition count to gather logs in parallel.
// Return 1 - 16 values depending on the expected log volume by the log filter.
TimePartitionCount(ctx context.Context) (int, error)
// Description returns the description for the Cloud Logging filter task.
Description() *ListLogEntriesTaskDescription
}
ListLogEntriesTaskSetting defines the settings for a Cloud Logging list log entries task.
type LocationFetcher ¶
type LocationFetcher interface {
FetchRegions(ctx context.Context, projectId string) ([]string, error)
}
func NewLocationFetcher ¶
func NewLocationFetcher(client *compute.RegionsClient, callOptionInjector *googlecloud.CallOptionInjector) LocationFetcher
type LogFetchProgress ¶
type LogFetchProgress struct {
// LogCount is the total number of logs fetched so far.
LogCount int
// Progress indicates the completion status, ranging from 0.0 to 1.0.
Progress float32
}
LogFetchProgress represents the progress of a log fetching operation.
type LogFetcher ¶
type LogFetcher interface {
FetchLogs(dest chan<- *loggingpb.LogEntry, ctx context.Context, filter string, container googlecloud.ResourceContainer, resourceContainers []string) error
}
LogFetcher is an interface for fetching logs from Cloud Logging with a given filter and sending them to a specified channel. The implementation must close the destination channel after the query is done.
func NewLogFetcher ¶
func NewLogFetcher(clientFactory *googlecloud.ClientFactory, callOptionInjector *googlecloud.CallOptionInjector, pageSize int32) LogFetcher
NewLogFetcher returns the instance of LogFetcher initialized with the given *googlecloud.ClientFactory.
type ProgressReportableLogFetcher ¶
type ProgressReportableLogFetcher interface {
// FetchLogsWithProgress fetches logs while periodically reporting its progress through a separate channel.
// It closes the progress channel upon completion and returns timestamp-sorted logs.
FetchLogsWithProgress(progress chan<- LogFetchProgress, ctx context.Context, beginTime, endTime time.Time, filterWithoutTimeRange string, container googlecloud.ResourceContainer, resourceContainers []string) ([]*log.Log, error)
}
type QueryResourceNames ¶
type QueryResourceNames struct {
QueryID string
DefaultResourceNames []string
CurrentResourceNames []string
}
QueryResourceNames holds the resource names for a specific query.
func (*QueryResourceNames) GetInputID ¶
func (q *QueryResourceNames) GetInputID() string
GetInputID returns the form input ID for the query.
type ResourceNamesInput ¶
type ResourceNamesInput struct {
// contains filtered or unexported fields
}
ResourceNamesInput is a container for resource names used in log queries.
func NewResourceNamesInput ¶
func NewResourceNamesInput() *ResourceNamesInput
NewResourceNamesInput creates a new ResourceNamesInput.
func (*ResourceNamesInput) GetResourceNamesForQuery ¶
func (r *ResourceNamesInput) GetResourceNamesForQuery(ctx context.Context, queryID string) *QueryResourceNames
GetResourceNamesForQuery returns the resource names for a given query ID.
func (*ResourceNamesInput) UpdateDefaultResourceNamesForQuery ¶
func (r *ResourceNamesInput) UpdateDefaultResourceNamesForQuery(queryID string, defaultResourceNames []string)
UpdateDefaultResourceNamesForQuery updates the default resource names for a given query ID.
type StandardProgressReportableLogFetcher ¶
type StandardProgressReportableLogFetcher struct {
// contains filtered or unexported fields
}
StandardProgressReportableLogFetcher is a decorator for a LogFetcher that adds the ability to report the progress of log fetching.
func NewStandardProgressReportableLogFetcher ¶
func NewStandardProgressReportableLogFetcher(fetcher LogFetcher, interval time.Duration) *StandardProgressReportableLogFetcher
NewProgressReportableLogFetcher creates a new instance of ProgressReportableLogFetcher.
func (*StandardProgressReportableLogFetcher) FetchLogsWithProgress ¶
func (s *StandardProgressReportableLogFetcher) FetchLogsWithProgress(dest chan<- *loggingpb.LogEntry, progress chan<- LogFetchProgress, ctx context.Context, beginTime, endTime time.Time, filterWithoutTimeRange string, container googlecloud.ResourceContainer, resourceContainers []string) error
FetchLogsWithProgress implements FetchLogsWithProgress.
type StructuredListLogEntriesTaskSetting ¶ added in v0.58.2
type StructuredListLogEntriesTaskSetting interface {
// TaskID returns the task ID for the structured list log entries task.
TaskID() taskid.TaskImplementationID[[]*log.Log]
// Dependencies returns the list of dependencies for the task.
Dependencies() []taskid.UntypedTaskReference
// DefaultResourceNames returns default resource names (e.g. ["projects/<project-id>"]).
DefaultResourceNames(ctx context.Context) ([]string, error)
// Queries returns the list of structured log queries for estimation and execution.
Queries(ctx context.Context) ([]*logestimator.StructuredLogQuery, error)
// TimePartitionCount returns the number of time partitions to gather logs in parallel.
TimePartitionCount(ctx context.Context) (int, error)
// QueryName returns human-readable name of the query.
QueryName() string
}
StructuredListLogEntriesTaskSetting defines the settings for a Cloud Logging task driven by StructuredLogQuery.
type TimePartitioningProgressReportableLogFetcher ¶
type TimePartitioningProgressReportableLogFetcher struct {
// contains filtered or unexported fields
}
func NewTimePartitioningProgressReportableLogFetcher ¶
func NewTimePartitioningProgressReportableLogFetcher(fetcher LogFetcher, interval time.Duration, partitionCount int, maxParallelism int) *TimePartitioningProgressReportableLogFetcher
func (*TimePartitioningProgressReportableLogFetcher) FetchLogsWithProgress ¶
func (t *TimePartitioningProgressReportableLogFetcher) FetchLogsWithProgress(progressChan chan<- LogFetchProgress, ctx context.Context, beginTime time.Time, endTime time.Time, filterWithoutTimeRange string, container googlecloud.ResourceContainer, resourceContainers []string) ([]*log.Log, error)
FetchLogsWithProgress implements ProgressReportableLogFetcher.
Source Files
¶
- contextkey.go
- extractor.go
- formpriority.go
- listlogentries_task.go
- locationfetcher.go
- log_ingester.go
- logfetcher.go
- operation_tracker.go
- progressreportablelogfetcher.go
- resourcename.go
- revision_state.go
- selectionpriority.go
- severity.go
- structured_listlogentries_task.go
- taskid.go
- tasklabel.go
- timeline.go
- timeline_type.go
- verb.go