Documentation
¶
Overview ¶
This package provides inspection tasks for Kubernetes node logs from Google Cloud Logging. The following is a Mermaid graph of the task dependencies within this package.
```mermaid graph TD
subgraph Inputs
InputProjectId
InputClusterName
InputNodeNameFilter
end
subgraph Log Fetching
ListLogEntries
end
subgraph Common Processing
LogSerializer
CommonFieldSetReader
end
subgraph Containerd Pipeline
ContainerdLogFilter
ContainerdLogGroup
ContainerdIDDiscovery
ContainerdLogToTimelineMapper
end
subgraph Kubelet Pipeline
KubeletLogFilter
KubeletLogGroup
KubeletLogToTimelineMapper
end
subgraph Other Pipeline
OtherLogFilter
OtherLogGroup
OtherLogToTimelineMapper
end
subgraph Finalization
Tail
end
%% Input Dependencies
InputProjectId --> ListLogEntries
InputClusterName --> ListLogEntries
InputNodeNameFilter --> ListLogEntries
%% Common Processing Dependencies
ListLogEntries --> LogSerializer
ListLogEntries --> CommonFieldSetReader
%% Containerd Pipeline Dependencies
CommonFieldSetReader --> ContainerdLogFilter
ContainerdLogFilter --> ContainerdLogGroup
ContainerdLogFilter --> ContainerdIDDiscovery
ContainerdIDDiscovery --> ContainerdLogToTimelineMapper
ContainerdLogGroup --> ContainerdLogToTimelineMapper
LogSerializer --> ContainerdLogToTimelineMapper
%% Kubelet Pipeline Dependencies
CommonFieldSetReader --> KubeletLogFilter
KubeletLogFilter --> KubeletLogGroup
ContainerdIDDiscovery --> KubeletLogToTimelineMapper
KubeletLogGroup --> KubeletLogToTimelineMapper
LogSerializer --> KubeletLogToTimelineMapper
%% Other Pipeline Dependencies
CommonFieldSetReader --> OtherLogFilter
OtherLogFilter --> OtherLogGroup
OtherLogGroup --> OtherLogToTimelineMapper
LogSerializer --> OtherLogToTimelineMapper
%% Finalization
ContainerdLogToTimelineMapper --> Tail
KubeletLogToTimelineMapper --> Tail
OtherLogToTimelineMapper --> Tail
```
Index ¶
- Constants
- Variables
- func GenerateK8sNodeLogQuery(cluster googlecloudk8scommon_contract.GoogleCloudClusterIdentity, ...) string
- func MustK8sNodeTimeline(ctx context.Context, clusterName string, nodeName string) *khifilev6.TimelinePath
- func MustK8sPodTimeline(ctx context.Context, clusterName string, namespace string, podName string) *khifilev6.TimelinePath
- func Register(registry coreinspection.InspectionTaskRegistry) error
- type K8sNodeLogIngester
Constants ¶
const ContainerdStartingMsg = "starting containerd"
const ContainerdTerminationMsg = "Stop CRI service"
Variables ¶
var ClusterIdentityAliasTask = coretask.NewAliasTask( googlecloudlogk8snode_contract.ClusterIdentityTaskID, googlecloudk8scommon_contract.ClusterIdentityTaskID.Ref(), )
var CommonFieldSetReaderTask = inspectiontaskbase.NewFieldSetReadTask(googlecloudlogk8snode_contract.CommonFieldsetReaderTaskID, googlecloudlogk8snode_contract.ListLogEntriesTaskID.Ref(), []log.FieldSetReader{ &googlecloudlogk8snode_contract.K8sNodeLogCommonFieldSetReader{ StructuredLogParser: logutil.NewMultiTextLogParser( logutil.NewJsonlTextParser(), logutil.NewKLogTextParser(true), logutil.NewLogfmtTextParser(), &logutil.FallbackRawTextLogParser{}, ), }, })
CommonFieldSetReaderTask parses the common fieldset used by GKE Node component logs.
var ContainerIDDiscoveryTask = commonlogk8saudit_contract.ContainerIDInventoryBuilder.DiscoveryTask(googlecloudlogk8snode_contract.ContainerIDDiscoveryTaskID, []taskid.UntypedTaskReference{ googlecloudlogk8snode_contract.ContainerdLogFilterTaskID.Ref(), }, func(ctx context.Context, taskMode inspectioncore_contract.InspectionTaskModeType, progress *inspectionmetadata.TaskProgressMetadata) (commonlogk8saudit_contract.ContainerIDToContainerIdentity, error) { if taskMode == inspectioncore_contract.TaskModeDryRun { return nil, nil } logs := coretask.GetTaskResult(ctx, googlecloudlogk8snode_contract.ContainerdLogFilterTaskID.Ref()) doneLogCount := atomic.Int32{} updator := progressutil.NewProgressUpdator(progress, time.Second, func(tp *inspectionmetadata.TaskProgressMetadata) { current := doneLogCount.Load() if len(logs) > 0 { tp.Percentage = float32(current) / float32(len(logs)) } tp.Message = fmt.Sprintf("%d/%d", current, len(logs)) }) updator.Start(ctx) defer updator.Done() result := commonlogk8saudit_contract.ContainerIDToContainerIdentity{} logChan := make(chan *log.Log) errGrp, childRoutineCtx := errgroup.WithContext(ctx) containerIdentitiesChan := make(chan *commonlogk8saudit_contract.ContainerIdentity, runtime.GOMAXPROCS(0)) for i := 0; i < runtime.GOMAXPROCS(0); i++ { errGrp.Go(func() error { for { select { case <-childRoutineCtx.Done(): return childRoutineCtx.Err() case l, ok := <-logChan: if !ok { return nil } processContainerIDDiscoveryForLog(ctx, l, containerIdentitiesChan) doneLogCount.Add(1) } } }) } consumerGrp, childConsumerRoutineCtx := errgroup.WithContext(ctx) consumerGrp.Go(func() error { for { select { case <-childConsumerRoutineCtx.Done(): return childConsumerRoutineCtx.Err() case c, ok := <-containerIdentitiesChan: if !ok { return nil } result[c.ContainerID] = c } } }) for _, l := range logs { logChan <- l } close(logChan) err := errGrp.Wait() close(containerIdentitiesChan) consumerErr := consumerGrp.Wait() if err != nil { return nil, err } if consumerErr != nil { return nil, consumerErr } return result, nil }, )
ContainerIDDiscoveryTask discovers mappings between container IDs and GKE pod containers.
var ContainerdLogFilterTask = newParserTypeFilterTask(googlecloudlogk8snode_contract.ContainerdLogFilterTaskID, googlecloudlogk8snode_contract.CommonFieldsetReaderTaskID.Ref(), googlecloudlogk8snode_contract.Containerd)
ContainerdLogFilterTask filters only containerd logs.
var ContainerdLogGroupTask = newNodeAndComponentNameGrouperTask(googlecloudlogk8snode_contract.ContainerdLogGroupTaskID, googlecloudlogk8snode_contract.ContainerdLogFilterTaskID.Ref())
ContainerdLogGroupTask groups containerd logs by node and component.
var ContainerdNodeLogLogToTimelineMapperTask = inspectiontaskbase.NewLogToTimelineMapperTaskV2( googlecloudlogk8snode_contract.ContainerdLogLogToTimelineMapperTaskID, &containerdNodeLogLogToTimelineMapperSetting{}, )
ContainerdNodeLogLogToTimelineMapperTask registers the mapper for containerd node component logs.
var KubeletLogFilterTask = newParserTypeFilterTask(googlecloudlogk8snode_contract.KubeletLogFilterTaskID, googlecloudlogk8snode_contract.CommonFieldsetReaderTaskID.Ref(), googlecloudlogk8snode_contract.Kubelet)
KubeletLogFilterTask filters only kubelet component logs.
var KubeletLogGroupTask = newNodeAndComponentNameGrouperTask(googlecloudlogk8snode_contract.KubeletLogGroupTaskID, googlecloudlogk8snode_contract.KubeletLogFilterTaskID.Ref())
KubeletLogGroupTask groups kubelet logs by node and component.
var KubeletLogLogToTimelineMapperTask = inspectiontaskbase.NewLogToTimelineMapperTaskV2( googlecloudlogk8snode_contract.KubeletLogLogToTimelineMapperTaskID, &kubeletNodeLogLogToTimelineMapperSetting{}, )
KubeletLogLogToTimelineMapperTask registers the mapper for kubelet component logs.
var ListLogEntriesTask = googlecloudcommon_contract.NewListLogEntriesTask(&k8snodeListLogEntriesTaskSetting{})
var LogIngesterTask = inspectiontaskbase.NewLogIngesterTaskV2( googlecloudlogk8snode_contract.LogIngesterTaskID, &K8sNodeLogIngester{}, )
LogIngesterTask registers the LogIngesterV2 for GKE Node logs.
var OtherLogFilterTask = newParserTypeFilterTask(googlecloudlogk8snode_contract.OtherLogFilterTaskID, googlecloudlogk8snode_contract.CommonFieldsetReaderTaskID.Ref(), googlecloudlogk8snode_contract.Other)
OtherLogFilterTask filters only the components logs that do not match kubelet or containerd.
var OtherLogGroupTask = newNodeAndComponentNameGrouperTask(googlecloudlogk8snode_contract.OtherLogGroupTaskID, googlecloudlogk8snode_contract.OtherLogFilterTaskID.Ref())
OtherLogGroupTask groups other logs by node and component.
var OtherLogLogToTimelineMapperTask = inspectiontaskbase.NewLogToTimelineMapperTaskV2(googlecloudlogk8snode_contract.OtherLogLogToTimelineMapperTaskID, &otherNodeLogLogToTimelineMapperSetting{ StartingMessagesByComponent: map[string]string{ "dockerd": "Starting up", "configure.sh": "Start to install kubernetes files", "configure-helper.sh": "Start to configure instance for kubernetes", }, TerminatingMessagesByComponent: map[string]string{ "dockerd": "Daemon shutdown complete", "configure.sh": "Done for installing kubernetes files", "configure-helper.sh": "Done for the configuration for kubernetes", }, })
OtherLogLogToTimelineMapperTask registers the mapper for other node component logs.
var PodSandboxIDDiscoveryTask = inspectiontaskbase.NewProgressReportableInspectionTask(googlecloudlogk8snode_contract.PodSandboxIDDiscoveryTaskID, []taskid.UntypedTaskReference{ googlecloudlogk8snode_contract.ContainerdLogFilterTaskID.Ref(), }, func(ctx context.Context, taskMode inspectioncore_contract.InspectionTaskModeType, progress *inspectionmetadata.TaskProgressMetadata) (patternfinder.PatternFinder[*googlecloudlogk8snode_contract.PodSandboxIDInfo], error) { if taskMode == inspectioncore_contract.TaskModeDryRun { return nil, nil } logs := coretask.GetTaskResult(ctx, googlecloudlogk8snode_contract.ContainerdLogFilterTaskID.Ref()) doneLogCount := atomic.Int32{} updator := progressutil.NewProgressUpdator(progress, time.Second, func(tp *inspectionmetadata.TaskProgressMetadata) { current := doneLogCount.Load() if len(logs) > 0 { tp.Percentage = float32(current) / float32(len(logs)) } tp.Message = fmt.Sprintf("%d/%d", current, len(logs)) }) updator.Start(ctx) defer updator.Done() logChan := make(chan *log.Log) errGrp, childCtx := errgroup.WithContext(ctx) podSandboxIDFinder := patternfinder.NewTriePatternFinder[*googlecloudlogk8snode_contract.PodSandboxIDInfo]() for i := 0; i < runtime.GOMAXPROCS(0); i++ { errGrp.Go(func() error { for { select { case <-childCtx.Done(): return childCtx.Err() case l, ok := <-logChan: if !ok { return nil } processPodSandboxIDDiscoveryForLog(ctx, l, podSandboxIDFinder) doneLogCount.Add(1) } } }) } for _, l := range logs { logChan <- l } close(logChan) errGrp.Wait() return podSandboxIDFinder, nil }, )
PodSandboxIDDiscoveryTask discovers mappings between pod sandbox IDs and GKE pods.
var TailTask = inspectiontaskbase.NewInspectionTask(googlecloudlogk8snode_contract.TailTaskID, []taskid.UntypedTaskReference{ googlecloudlogk8snode_contract.ContainerdLogLogToTimelineMapperTaskID.Ref(), googlecloudlogk8snode_contract.KubeletLogLogToTimelineMapperTaskID.Ref(), googlecloudlogk8snode_contract.OtherLogLogToTimelineMapperTaskID.Ref(), googlecloudlogk8snode_contract.ContainerIDDiscoveryTaskID.Ref(), }, func(ctx context.Context, taskMode inspectioncore_contract.InspectionTaskModeType) (struct{}, error) { return struct{}{}, nil }, inspectioncore_contract.FeatureTaskLabelV2( "Kubernetes Node Logs", "Gather logs from Kubernetes node components (e.g., Docker, containerd, or Kubelet) to troubleshoot node-level issues. Note: The log volume can be very large if the cluster contains many nodes.", 3000, false, ), )
TailTask is a nop task that depends on all node component mappers and other child tasks to group them.
Functions ¶
func GenerateK8sNodeLogQuery ¶
func GenerateK8sNodeLogQuery(cluster googlecloudk8scommon_contract.GoogleCloudClusterIdentity, nodeNameSubstrings []string) string
GenerateK8sNodeLogQuery generates a query for GKE node logs.
func MustK8sNodeTimeline ¶ added in v0.56.0
func MustK8sNodeTimeline(ctx context.Context, clusterName string, nodeName string) *khifilev6.TimelinePath
MustK8sNodeTimeline returns the timeline path for the Kubernetes Node resource layer.
func MustK8sPodTimeline ¶ added in v0.56.0
func MustK8sPodTimeline(ctx context.Context, clusterName string, namespace string, podName string) *khifilev6.TimelinePath
MustK8sPodTimeline returns the timeline path for a Kubernetes Pod resource layer.
func Register ¶
func Register(registry coreinspection.InspectionTaskRegistry) error
Register registers all googlecloudlogk8snode inspection tasks to the registry.
Types ¶
type K8sNodeLogIngester ¶ added in v0.56.0
type K8sNodeLogIngester struct{}
K8sNodeLogIngester implements LogIngesterV2 for GKE Node component logs.
func (*K8sNodeLogIngester) Dependencies ¶ added in v0.56.0
func (i *K8sNodeLogIngester) Dependencies() []taskid.UntypedTaskReference
Dependencies returns the dependencies of the log ingester.
func (*K8sNodeLogIngester) ProcessLog ¶ added in v0.56.0
func (i *K8sNodeLogIngester) ProcessLog(ctx context.Context, l *log.Log) (*khifilev6.LogChangeSet, error)
ProcessLog populates the LogChangeSet for GKE Node logs.
func (*K8sNodeLogIngester) RawLogTask ¶ added in v0.56.0
func (i *K8sNodeLogIngester) RawLogTask() taskid.TaskReference[[]*log.Log]
RawLogTask returns the raw log provider task.