csm_impl

package
v0.59.0 Latest Latest
Warning

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

Go to latest
Published: Sep 17, 2026 License: Apache-2.0 Imports: 24 Imported by: 0

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

Constants

This section is empty.

Variables

View Source
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).

View Source
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.

CSMTrafficDirectorLogIngesterTask is a task that ingests CSM Traffic Director logs.

View Source
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.

View Source
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()
View Source
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()
View Source
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.

LogIngesterTask is the task that executes CSMTrafficLogLogIngester.

View Source
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

Dependencies implements gcpcommon.StructuredListLogEntriesTaskSetting.

func (*CSMTrafficDirectorListLogEntryTaskSetting) Queries

Queries implements gcpcommon.StructuredListLogEntriesTaskSetting.

func (*CSMTrafficDirectorListLogEntryTaskSetting) QueryName

QueryName implements gcpcommon.StructuredListLogEntriesTaskSetting.

func (*CSMTrafficDirectorListLogEntryTaskSetting) TaskID

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

Dependencies returns additional task dependencies.

func (*CSMTrafficDirectorLogToTimelineMapper) GroupedLogTask

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

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

Dependencies implements gcpcommon.StructuredListLogEntriesTaskSetting.

func (*CSMTrafficLogListLogEntryTaskSetting) Queries

Queries implements gcpcommon.StructuredListLogEntriesTaskSetting.

func (*CSMTrafficLogListLogEntryTaskSetting) QueryName

QueryName implements gcpcommon.StructuredListLogEntriesTaskSetting.

func (*CSMTrafficLogListLogEntryTaskSetting) TaskID

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

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

Dependencies returns additional task dependencies.

func (*CSMTrafficLogLogToTimelineMapper) GroupedLogTask

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.

Jump to

Keyboard shortcuts

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