Documentation
¶
Overview ¶
Copyright 2026 Google LLC
Licensed under the Apache License, Version 2.0 (the "License"); you may not use this file except in compliance with the License. You may obtain a copy of the License at
http://www.apache.org/licenses/LICENSE-2.0
Unless required by applicable law or agreed to in writing, software distributed under the License is distributed on an "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the License for the specific language governing permissions and limitations under the License.
Index ¶
- Variables
- func GenerateCSMTrafficDirectorStructuredQuery(fleetProjectID string, clusterIdentifiers []string, isDryRun bool) *logestimator.StructuredLogQuery
- func GenerateCSMTrafficLogsStructuredQuery(cluster k8scommon.GoogleCloudClusterIdentity, ...) *logestimator.StructuredLogQuery
- func Register(registry coreinspection.InspectionTaskRegistry) error
- type CSMTrafficDirectorListLogEntryTaskSetting
- func (s *CSMTrafficDirectorListLogEntryTaskSetting) DefaultResourceNames(ctx context.Context) ([]string, error)
- func (s *CSMTrafficDirectorListLogEntryTaskSetting) Dependencies() []coretask.Dependency
- func (s *CSMTrafficDirectorListLogEntryTaskSetting) Queries(ctx context.Context) ([]*logestimator.StructuredLogQuery, error)
- func (s *CSMTrafficDirectorListLogEntryTaskSetting) QueryName() string
- func (s *CSMTrafficDirectorListLogEntryTaskSetting) TaskID() taskid.TaskImplementationID[[]*log.Log]
- func (s *CSMTrafficDirectorListLogEntryTaskSetting) TimePartitionCount(ctx context.Context) (int, error)
- type CSMTrafficDirectorLogToTimelineMapper
- func (m *CSMTrafficDirectorLogToTimelineMapper) Dependencies() []coretask.Dependency
- func (m *CSMTrafficDirectorLogToTimelineMapper) GroupedLogTask() taskid.TaskReference[inspectiontaskbase.LogGroupMap]
- func (m *CSMTrafficDirectorLogToTimelineMapper) LogIngesterTask() taskid.TaskReference[struct{}]
- func (m *CSMTrafficDirectorLogToTimelineMapper) ProcessLogByGroup(ctx context.Context, l *log.Log, tracker *gcpcommon.GCPOperationTracker) (*khifilev6.TimelineChangeSet, *gcpcommon.GCPOperationTracker, error)
- type CSMTrafficLogListLogEntryTaskSetting
- func (c *CSMTrafficLogListLogEntryTaskSetting) DefaultResourceNames(ctx context.Context) ([]string, error)
- func (c *CSMTrafficLogListLogEntryTaskSetting) Dependencies() []coretask.Dependency
- func (c *CSMTrafficLogListLogEntryTaskSetting) Queries(ctx context.Context) ([]*logestimator.StructuredLogQuery, error)
- func (c *CSMTrafficLogListLogEntryTaskSetting) QueryName() string
- func (c *CSMTrafficLogListLogEntryTaskSetting) TaskID() taskid.TaskImplementationID[[]*log.Log]
- func (c *CSMTrafficLogListLogEntryTaskSetting) TimePartitionCount(ctx context.Context) (int, error)
- type CSMTrafficLogLogIngester
- type CSMTrafficLogLogToTimelineMapper
- func (m *CSMTrafficLogLogToTimelineMapper) Dependencies() []coretask.Dependency
- func (m *CSMTrafficLogLogToTimelineMapper) GroupedLogTask() taskid.TaskReference[inspectiontaskbase.LogGroupMap]
- func (m *CSMTrafficLogLogToTimelineMapper) LogIngesterTask() taskid.TaskReference[struct{}]
- func (m *CSMTrafficLogLogToTimelineMapper) ProcessLogByGroup(ctx context.Context, l *log.Log, _ struct{}) (*khifilev6.TimelineChangeSet, struct{}, error)
Constants ¶
This section is empty.
Variables ¶
var CSMClusterIdentifierTask = coretask.NewTask( csm.CSMClusterIdentifierTaskID, []coretask.Dependency{k8scommon.NEGToBackendServiceInventoryTaskID.Ref()}, func(ctx context.Context) ([]string, error) { inventory := coretask.GetTaskResult(ctx, k8scommon.NEGToBackendServiceInventoryTaskID.Ref()) uniqueIds := make(map[string]struct{}) for _, bsName := range inventory { if strings.HasPrefix(bsName, "gsmrsvd-") { parts := strings.Split(bsName, "-") if len(parts) >= 3 { uniqueIds[parts[1]] = struct{}{} } } } result := make([]string, 0, len(uniqueIds)) for id := range uniqueIds { result = append(result, id) } return result, nil }, )
CSMClusterIdentifierTask extracts the unique cluster identifier(s) from BackendService names. CSM BackendService names follow the pattern: gsmrsvd-(cluster-identifier)-(neg-id).
var CSMTrafficDirectorLogGrouperTask = inspectiontaskbase.NewLogGrouperTask( csm.CSMTrafficDirectorLogGrouperTaskID, csm.ListCSMTrafficDirectorLogEntriesTaskID.Ref(), func(ctx context.Context, l *log.Log) string { audit, err := gcpcommon.ExtractGCPAuditLog(l.NodeReader) if err != nil { return "unknown" } return audit.ResourceName }, )
CSMTrafficDirectorLogGrouperTask is a task that groups CSM Traffic Director logs by their resource name.
var CSMTrafficDirectorLogIngesterTask = gcpcommon.NewGCPOperationLogIngesterTask( csm.CSMTrafficDirectorLogIngesterTaskID, csm.ListCSMTrafficDirectorLogEntriesTaskID.Ref(), csm.LogTypeCSMTrafficLog, )
CSMTrafficDirectorLogIngesterTask is a task that ingests CSM Traffic Director logs.
var CSMTrafficDirectorLogToTimelineMapperTask = inspectiontaskbase.NewLogToTimelineMapperTask( csm.CSMTrafficDirectorLogToTimelineMapperTaskID, &CSMTrafficDirectorLogToTimelineMapper{}, inspectioncore.FeatureTaskLabel( "CSM Resource Audit Logs", "Gather audit logs for Traffic Director resources created by the TD-based CSM to map them to timelines alongside associated Kubernetes resource logs.", 10100, false, ), )
CSMTrafficDirectorLogToTimelineMapperTask maps CSM Traffic Director logs to timelines.
var ClusterIdentityAliasTask = coretask.NewAliasTask( csm.ClusterIdentityTaskID, k8scommon.ClusterIdentityTaskID.Ref(), )
var InputCSMResponseFlagsTask = formtask.NewSetFormTaskBuilder(csm.InputCSMResponseFlagsTaskID, priorityForCSMGroup+1000, "Envoy response flags"). WithDefaultValueConstant([]string{"@any", "-OK"}, true). WithAllowAddAll(false). WithAllowRemoveAll(false). WithAllowCustomValue(true). WithDescription("Response flags used for filtering CSM traffic logs. Note '-' in response flags is corresponded to 'OK' in this form."). WithOptionsFunc(func(ctx context.Context, previousValues []string) ([]inspectionmetadata.SetParameterFormFieldOptionItem, error) { result := []inspectionmetadata.SetParameterFormFieldOptionItem{ {ID: "@any", Description: "[Alias] Matches any response flag"}, } ids := make([]string, 0, len(logutil.EnvoyResponseFlagDescriptions)) for flag := range logutil.EnvoyResponseFlagDescriptions { ids = append(ids, string(flag)) } sort.Strings(ids) for _, id := range ids { message := logutil.EnvoyResponseFlagDescriptions[logutil.EnvoyResponseFlag(id)] if id == "-" { id = "OK" message = "It's '-' in the response flag field because '-' means subtracting operator in this form." } result = append(result, inspectionmetadata.SetParameterFormFieldOptionItem{ID: id, Description: message}) } return result, nil }). WithValidator(func(ctx context.Context, value []string) (string, error) { strFilter := strings.Join(value, " ") result, err := gcpqueryutil.ParseSetFilter(strFilter, inputCSMAliasMap, true, true, true) if err != nil { return "", err } if result.ValidationError == "" { err = verifyResponseFlags(convertInputOnlyResponseFlagToActualFlag(result.Additives)) if err != nil { return err.Error(), nil } err = verifyResponseFlags(convertInputOnlyResponseFlagToActualFlag(result.Subtractives)) if err != nil { return err.Error(), nil } } return result.ValidationError, nil }). WithConverter(func(ctx context.Context, value []string) (*gcpqueryutil.SetFilterParseResult, error) { strFilter := strings.Join(value, " ") result, err := gcpqueryutil.ParseSetFilter(strFilter, inputCSMAliasMap, true, true, true) if err != nil { return nil, err } result.Additives = convertInputOnlyResponseFlagToActualFlag(result.Additives) result.Subtractives = convertInputOnlyResponseFlagToActualFlag(result.Subtractives) return result, nil }). Build()
var InputFleetProjectIDTask = formtask.NewTextFormTaskBuilder(csm.InputFleetProjectIDTaskID, priorityForCSMGroup+500, "Fleet project ID"). WithDependencies([]coretask.Dependency{ csm.ClusterIdentityTaskID.Ref(), }). WithDescription("The project ID where the Fleet is hosted and CSM control plane logs are stored. Default is the cluster's project ID."). WithDefaultValueFunc(func(ctx context.Context, previousValues []string) (string, error) { if len(previousValues) > 0 { return previousValues[0], nil } cluster := coretask.GetTaskResult(ctx, csm.ClusterIdentityTaskID.Ref()) return cluster.ProjectID, nil }). Build()
var ListCSMTrafficDirectorLogEntriesTask = gcpcommon.NewStructuredListLogEntriesTask(&CSMTrafficDirectorListLogEntryTaskSetting{})
var ListLogEntriesTask = gcpcommon.NewStructuredListLogEntriesTask(&CSMTrafficLogListLogEntryTaskSetting{})
var LogGrouperTask = inspectiontaskbase.NewLogGrouperTask(csm.LogGrouperTaskID, csm.ListLogEntriesTaskID.Ref(), func(ctx context.Context, l *log.Log) string { istioAccessLogFieldSet, err := csm.ExtractIstioAccessLog(l.NodeReader) if err != nil { return "unknown" } return fmt.Sprintf("%s-%s", istioAccessLogFieldSet.ReporterPodNamespace, istioAccessLogFieldSet.ReporterPodName) }, )
LogGrouperTask groups CSM traffic logs by their reporter pod.
var LogIngesterTask = inspectiontaskbase.NewLogIngesterTask( csm.LogIngesterTaskID, &CSMTrafficLogLogIngester{}, )
LogIngesterTask is the task that executes CSMTrafficLogLogIngester.
var LogToTimelineMapperTask = inspectiontaskbase.NewLogToTimelineMapperTask( csm.LogToTimelineMapperTaskID, &CSMTrafficLogLogToTimelineMapper{}, inspectioncore.FeatureTaskLabel( "CSM Traffic Logs", "Gather CSM traffic logs to visualize network traffic flows and latency under client or server Pod timelines.", 10000, false, ), )
LogToTimelineMapperTask maps CSM traffic logs to timelines.
Functions ¶
func GenerateCSMTrafficDirectorStructuredQuery ¶
func GenerateCSMTrafficDirectorStructuredQuery(fleetProjectID string, clusterIdentifiers []string, isDryRun bool) *logestimator.StructuredLogQuery
GenerateCSMTrafficDirectorStructuredQuery generates a structured query for CSM Traffic Director logs.
func GenerateCSMTrafficLogsStructuredQuery ¶
func GenerateCSMTrafficLogsStructuredQuery(cluster k8scommon.GoogleCloudClusterIdentity, responseFlagsSetFilter *gcpqueryutil.SetFilterParseResult, namespaceSetFilter *gcpqueryutil.SetFilterParseResult) *logestimator.StructuredLogQuery
GenerateCSMTrafficLogsStructuredQuery generates a structured query for CSM Traffic logs.
func Register ¶
func Register(registry coreinspection.InspectionTaskRegistry) error
graph TD subgraph "CSM Traffic Log" direction LR InputCSMResponseFlagsTask(Input CSM Response Flags) ListLogEntriesTask(List Log Entries) LogIngesterTask(Log Serializer) LogGrouperTask(Log Grouper) LogToTimelineMapperTask(TimelineMapper) ListLogEntriesTask --> LogGrouperTask ListLogEntriesTask --> LogIngesterTask LogGrouperTask --> LogToTimelineMapperTask LogIngesterTask --> LogToTimelineMapperTask InputCSMResponseFlagsTask --> ListLogEntriesTask end
Register registers all googlecloudlogcsm inspection tasks to the registry.
Types ¶
type CSMTrafficDirectorListLogEntryTaskSetting ¶
type CSMTrafficDirectorListLogEntryTaskSetting struct{}
func (*CSMTrafficDirectorListLogEntryTaskSetting) DefaultResourceNames ¶
func (s *CSMTrafficDirectorListLogEntryTaskSetting) DefaultResourceNames(ctx context.Context) ([]string, error)
DefaultResourceNames implements gcpcommon.StructuredListLogEntriesTaskSetting.
func (*CSMTrafficDirectorListLogEntryTaskSetting) Dependencies ¶
func (s *CSMTrafficDirectorListLogEntryTaskSetting) Dependencies() []coretask.Dependency
Dependencies implements gcpcommon.StructuredListLogEntriesTaskSetting.
func (*CSMTrafficDirectorListLogEntryTaskSetting) Queries ¶
func (s *CSMTrafficDirectorListLogEntryTaskSetting) Queries(ctx context.Context) ([]*logestimator.StructuredLogQuery, error)
Queries implements gcpcommon.StructuredListLogEntriesTaskSetting.
func (*CSMTrafficDirectorListLogEntryTaskSetting) QueryName ¶
func (s *CSMTrafficDirectorListLogEntryTaskSetting) QueryName() string
QueryName implements gcpcommon.StructuredListLogEntriesTaskSetting.
func (*CSMTrafficDirectorListLogEntryTaskSetting) TaskID ¶
func (s *CSMTrafficDirectorListLogEntryTaskSetting) TaskID() taskid.TaskImplementationID[[]*log.Log]
TaskID implements gcpcommon.StructuredListLogEntriesTaskSetting.
func (*CSMTrafficDirectorListLogEntryTaskSetting) TimePartitionCount ¶
func (s *CSMTrafficDirectorListLogEntryTaskSetting) TimePartitionCount(ctx context.Context) (int, error)
TimePartitionCount implements gcpcommon.StructuredListLogEntriesTaskSetting.
type CSMTrafficDirectorLogToTimelineMapper ¶
type CSMTrafficDirectorLogToTimelineMapper struct {
inspectiontaskbase.SinglePassMapperBase[*gcpcommon.GCPOperationTracker]
}
CSMTrafficDirectorLogToTimelineMapper maps CSM Traffic Director logs to resource timelines.
func (*CSMTrafficDirectorLogToTimelineMapper) Dependencies ¶
func (m *CSMTrafficDirectorLogToTimelineMapper) Dependencies() []coretask.Dependency
Dependencies returns additional task dependencies.
func (*CSMTrafficDirectorLogToTimelineMapper) GroupedLogTask ¶
func (m *CSMTrafficDirectorLogToTimelineMapper) GroupedLogTask() taskid.TaskReference[inspectiontaskbase.LogGroupMap]
GroupedLogTask returns a reference to the task that provides the grouped logs.
func (*CSMTrafficDirectorLogToTimelineMapper) LogIngesterTask ¶
func (m *CSMTrafficDirectorLogToTimelineMapper) LogIngesterTask() taskid.TaskReference[struct{}]
LogIngesterTask returns a reference to the task that provides ingested logs.
func (*CSMTrafficDirectorLogToTimelineMapper) ProcessLogByGroup ¶
func (m *CSMTrafficDirectorLogToTimelineMapper) ProcessLogByGroup(ctx context.Context, l *log.Log, tracker *gcpcommon.GCPOperationTracker) (*khifilev6.TimelineChangeSet, *gcpcommon.GCPOperationTracker, error)
ProcessLogByGroup maps each log inside a group to one or more timeline events or revisions.
type CSMTrafficLogListLogEntryTaskSetting ¶
type CSMTrafficLogListLogEntryTaskSetting struct{}
func (*CSMTrafficLogListLogEntryTaskSetting) DefaultResourceNames ¶
func (c *CSMTrafficLogListLogEntryTaskSetting) DefaultResourceNames(ctx context.Context) ([]string, error)
DefaultResourceNames implements gcpcommon.StructuredListLogEntriesTaskSetting.
func (*CSMTrafficLogListLogEntryTaskSetting) Dependencies ¶
func (c *CSMTrafficLogListLogEntryTaskSetting) Dependencies() []coretask.Dependency
Dependencies implements gcpcommon.StructuredListLogEntriesTaskSetting.
func (*CSMTrafficLogListLogEntryTaskSetting) Queries ¶
func (c *CSMTrafficLogListLogEntryTaskSetting) Queries(ctx context.Context) ([]*logestimator.StructuredLogQuery, error)
Queries implements gcpcommon.StructuredListLogEntriesTaskSetting.
func (*CSMTrafficLogListLogEntryTaskSetting) QueryName ¶
func (c *CSMTrafficLogListLogEntryTaskSetting) QueryName() string
QueryName implements gcpcommon.StructuredListLogEntriesTaskSetting.
func (*CSMTrafficLogListLogEntryTaskSetting) TaskID ¶
func (c *CSMTrafficLogListLogEntryTaskSetting) TaskID() taskid.TaskImplementationID[[]*log.Log]
TaskID implements gcpcommon.StructuredListLogEntriesTaskSetting.
func (*CSMTrafficLogListLogEntryTaskSetting) TimePartitionCount ¶
func (c *CSMTrafficLogListLogEntryTaskSetting) TimePartitionCount(ctx context.Context) (int, error)
TimePartitionCount implements gcpcommon.StructuredListLogEntriesTaskSetting.
type CSMTrafficLogLogIngester ¶
type CSMTrafficLogLogIngester struct{}
CSMTrafficLogLogIngester ingests CSM traffic logs.
func (*CSMTrafficLogLogIngester) Dependencies ¶
func (i *CSMTrafficLogLogIngester) Dependencies() []coretask.Dependency
Dependencies returns the task dependencies.
func (*CSMTrafficLogLogIngester) ProcessLog ¶
func (i *CSMTrafficLogLogIngester) ProcessLog(ctx context.Context, l *log.Log) (*khifilev6.LogChangeSet, error)
ProcessLog parses raw log entry and populates the LogChangeSet.
func (*CSMTrafficLogLogIngester) RawLogTask ¶
func (i *CSMTrafficLogLogIngester) RawLogTask() taskid.TaskReference[[]*log.Log]
RawLogTask returns the task reference that provides raw logs.
type CSMTrafficLogLogToTimelineMapper ¶
type CSMTrafficLogLogToTimelineMapper struct {
inspectiontaskbase.StatelessMapperBase
}
CSMTrafficLogLogToTimelineMapper maps CSM traffic logs to resource timelines.
func (*CSMTrafficLogLogToTimelineMapper) Dependencies ¶
func (m *CSMTrafficLogLogToTimelineMapper) Dependencies() []coretask.Dependency
Dependencies returns additional task dependencies.
func (*CSMTrafficLogLogToTimelineMapper) GroupedLogTask ¶
func (m *CSMTrafficLogLogToTimelineMapper) GroupedLogTask() taskid.TaskReference[inspectiontaskbase.LogGroupMap]
GroupedLogTask returns a reference to the task that provides the grouped logs.
func (*CSMTrafficLogLogToTimelineMapper) LogIngesterTask ¶
func (m *CSMTrafficLogLogToTimelineMapper) LogIngesterTask() taskid.TaskReference[struct{}]
LogIngesterTask returns a reference to the task that provides ingested logs.
func (*CSMTrafficLogLogToTimelineMapper) ProcessLogByGroup ¶
func (m *CSMTrafficLogLogToTimelineMapper) ProcessLogByGroup(ctx context.Context, l *log.Log, _ struct{}) (*khifilev6.TimelineChangeSet, struct{}, error)
ProcessLogByGroup maps each log inside a group to one or more timeline events.