k8snode_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: 21 Imported by: 0

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

%% Containerd Pipeline Dependencies
ListLogEntries --> ContainerdLogFilter
ContainerdLogFilter --> ContainerdLogGroup
ContainerdLogFilter --> ContainerdIDDiscovery
ContainerdIDDiscovery --> ContainerdLogToTimelineMapper
ContainerdLogGroup --> ContainerdLogToTimelineMapper
LogSerializer --> ContainerdLogToTimelineMapper

%% Kubelet Pipeline Dependencies
ListLogEntries --> KubeletLogFilter
KubeletLogFilter --> KubeletLogGroup
ContainerdIDDiscovery --> KubeletLogToTimelineMapper
KubeletLogGroup --> KubeletLogToTimelineMapper
LogSerializer --> KubeletLogToTimelineMapper

%% Other Pipeline Dependencies
ListLogEntries --> OtherLogFilter
OtherLogFilter --> OtherLogGroup
OtherLogGroup --> OtherLogToTimelineMapper
LogSerializer --> OtherLogToTimelineMapper

%% Finalization
ContainerdLogToTimelineMapper --> Tail
KubeletLogToTimelineMapper --> Tail
OtherLogToTimelineMapper --> Tail

```

Index

Constants

View Source
const ContainerdStartingMsg = "starting containerd"
View Source
const ContainerdTerminationMsg = "Stop CRI service"

Variables

View Source
var ContainerIDDiscoveryTask = inspectiontaskbase.NewInspectionTask(k8snode.ContainerIDDiscoveryTaskID,
	[]coretask.Dependency{
		k8snode.ContainerdLogFilterTaskID.Ref(),
	},
	func(ctx context.Context, taskMode inspectioncore.InspectionTaskModeType) (k8saudit.ContainerIDToContainerIdentity, error) {
		if taskMode == inspectioncore.TaskModeDryRun {
			return nil, nil
		}

		logs := coretask.GetTaskResult(ctx, k8snode.ContainerdLogFilterTaskID.Ref())

		tracker := progress.NewTracker(ctx, len(logs), progress.WithUnit("logs"))
		defer tracker.Done()

		result := k8saudit.ContainerIDToContainerIdentity{}
		logChan := make(chan *log.Log)
		errGrp, childRoutineCtx := errgroup.WithContext(ctx)
		containerIdentitiesChan := make(chan *k8saudit.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)
						tracker.Inc()
					}
				}
			})
		}
		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
	},
	coretask.ProvidesTag(k8saudit.TagContainerIDDiscovery),
)

ContainerIDDiscoveryTask discovers mappings between container IDs and GKE pod containers.

View Source
var ContainerdLogFilterTask = newParserTypeFilterTask(k8snode.ContainerdLogFilterTaskID, k8snode.ListLogEntriesTaskID.Ref(), k8snode.Containerd)

ContainerdLogFilterTask filters only containerd logs.

View Source
var ContainerdLogGroupTask = newNodeAndComponentNameGrouperTask(k8snode.ContainerdLogGroupTaskID, k8snode.ContainerdLogFilterTaskID.Ref())

ContainerdLogGroupTask groups containerd logs by node and component.

View Source
var ContainerdNodeLogLogToTimelineMapperTask = inspectiontaskbase.NewLogToTimelineMapperTask(
	k8snode.ContainerdLogLogToTimelineMapperTaskID,
	&containerdNodeLogLogToTimelineMapperSetting{},
)

ContainerdNodeLogLogToTimelineMapperTask registers the mapper for containerd node component logs.

View Source
var KubeletLogFilterTask = newParserTypeFilterTask(k8snode.KubeletLogFilterTaskID, k8snode.ListLogEntriesTaskID.Ref(), k8snode.Kubelet)

KubeletLogFilterTask filters only kubelet component logs.

View Source
var KubeletLogGroupTask = newNodeAndComponentNameGrouperTask(k8snode.KubeletLogGroupTaskID, k8snode.KubeletLogFilterTaskID.Ref())

KubeletLogGroupTask groups kubelet logs by node and component.

View Source
var KubeletLogLogToTimelineMapperTask = inspectiontaskbase.NewLogToTimelineMapperTask(
	k8snode.KubeletLogLogToTimelineMapperTaskID,
	&kubeletNodeLogLogToTimelineMapperSetting{},
)

KubeletLogLogToTimelineMapperTask registers the mapper for kubelet component logs.

View Source
var ListLogEntriesTask = gcpcommon.NewStructuredListLogEntriesTask(&k8snodeListLogEntriesTaskSetting{})

LogIngesterTask registers the LogIngester for GKE Node logs.

View Source
var NodeNameDiscoveryTask = inspectiontaskbase.NewInspectionTask(
	k8snode.NodeNameDiscoveryTaskID,
	[]coretask.Dependency{
		k8snode.ListLogEntriesTaskID.Ref(),
	},
	func(ctx context.Context, taskMode inspectioncore.InspectionTaskModeType) ([]string, error) {
		if taskMode == inspectioncore.TaskModeDryRun {
			return nil, nil
		}

		foundNodeNames := map[string]struct{}{}
		logs := coretask.GetTaskResult(ctx, k8snode.ListLogEntriesTaskID.Ref())
		for _, l := range logs {
			fs, err := k8snode.ExtractK8sNodeLogCommon(l.NodeReader, nil)
			if err == nil && fs.NodeName != "" {
				foundNodeNames[fs.NodeName] = struct{}{}
			}
		}

		var result []string
		for k := range foundNodeNames {
			result = append(result, k)
		}
		return result, nil
	},
	coretask.ProvidesTag(k8saudit.TagNodeNameDiscovery),
)

NodeNameDiscoveryTask extracts node names from Kubernetes Node component logs and registers them to NodeNameInventoryTask.

View Source
var OtherLogFilterTask = newParserTypeFilterTask(k8snode.OtherLogFilterTaskID, k8snode.ListLogEntriesTaskID.Ref(), k8snode.Other)

OtherLogFilterTask filters only the components logs that do not match kubelet or containerd.

View Source
var OtherLogGroupTask = newNodeAndComponentNameGrouperTask(k8snode.OtherLogGroupTaskID, k8snode.OtherLogFilterTaskID.Ref())

OtherLogGroupTask groups other logs by node and component.

View Source
var OtherLogLogToTimelineMapperTask = inspectiontaskbase.NewLogToTimelineMapperTask(k8snode.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.

View Source
var PodSandboxIDDiscoveryTask = inspectiontaskbase.NewInspectionTask(k8snode.PodSandboxIDDiscoveryTaskID,
	[]coretask.Dependency{
		k8snode.ContainerdLogFilterTaskID.Ref(),
	},
	func(ctx context.Context, taskMode inspectioncore.InspectionTaskModeType) (patternfinder.PatternFinder[*k8snode.PodSandboxIDInfo], error) {
		if taskMode == inspectioncore.TaskModeDryRun {
			return nil, nil
		}
		logs := coretask.GetTaskResult(ctx, k8snode.ContainerdLogFilterTaskID.Ref())

		tracker := progress.NewTracker(ctx, len(logs), progress.WithUnit("logs"))
		defer tracker.Done()

		logChan := make(chan *log.Log)
		errGrp, childCtx := errgroup.WithContext(ctx)
		podSandboxIDFinder := patternfinder.NewRadixPatternFinder[*k8snode.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)
						tracker.Inc()
					}
				}
			})
		}

		for _, l := range logs {
			logChan <- l
		}
		close(logChan)
		errGrp.Wait()

		return podSandboxIDFinder, nil
	},
)

PodSandboxIDDiscoveryTask discovers mappings between pod sandbox IDs and GKE pods.

View Source
var TailTask = coretask.NewTailTask(
	k8snode.TailTaskID,
	[]coretask.Dependency{
		k8snode.ContainerdLogLogToTimelineMapperTaskID.Ref(),
		k8snode.KubeletLogLogToTimelineMapperTaskID.Ref(),
		k8snode.OtherLogLogToTimelineMapperTaskID.Ref(),

		k8snode.ContainerIDDiscoveryTaskID.Ref(),
		k8snode.NodeNameDiscoveryTaskID.Ref(),
	},
	inspectioncore.FeatureTaskLabel(
		"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 GenerateK8sNodeStructuredQuery

func GenerateK8sNodeStructuredQuery(cluster k8scommon.GoogleCloudClusterIdentity, nodeNameSubstrings []string) *logestimator.StructuredLogQuery

GenerateK8sNodeStructuredQuery generates a structured query for GKE node logs.

func MustK8sNodeTimeline

func MustK8sNodeTimeline(ctx context.Context, clusterName string, nodeName string) *khifilev6.TimelinePath

MustK8sNodeTimeline returns the timeline path for the Kubernetes Node resource layer.

func MustK8sPodTimeline

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

type K8sNodeLogIngester struct{}

K8sNodeLogIngester implements LogIngester for GKE Node component logs.

func (*K8sNodeLogIngester) Dependencies

func (i *K8sNodeLogIngester) Dependencies() []coretask.Dependency

Dependencies returns the dependencies of the log ingester.

func (*K8sNodeLogIngester) ProcessLog

func (i *K8sNodeLogIngester) ProcessLog(ctx context.Context, l *log.Log) (*khifilev6.LogChangeSet, error)

ProcessLog populates the LogChangeSet for GKE Node logs.

func (*K8sNodeLogIngester) RawLogTask

func (i *K8sNodeLogIngester) RawLogTask() taskid.TaskReference[[]*log.Log]

RawLogTask returns the raw log provider task.

Jump to

Keyboard shortcuts

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