diff options
| author | fisherxu <xufei40@huawei.com> | 2020-03-28 18:41:25 +0800 |
|---|---|---|
| committer | fisherxu <xufei40@huawei.com> | 2020-04-17 17:50:46 +0800 |
| commit | 062b3ddd6b70bad35057359a273d443c5d972d72 (patch) | |
| tree | b8ab42d2eec8a736349a2fc7460176b3906f7fce /edge | |
| parent | add vendor (diff) | |
| download | kubeedge-062b3ddd6b70bad35057359a273d443c5d972d72.tar.gz | |
add cadvisor
Diffstat (limited to 'edge')
| -rw-r--r-- | edge/pkg/edged/cadvisor/cadvisor_linux.go | 98 | ||||
| -rw-r--r-- | edge/pkg/edged/containers/helpers.go | 19 | ||||
| -rw-r--r-- | edge/pkg/edged/edged.go | 136 | ||||
| -rw-r--r-- | edge/pkg/edged/edged_getters.go | 90 | ||||
| -rw-r--r-- | edge/pkg/edged/edged_pods.go | 169 | ||||
| -rw-r--r-- | edge/pkg/edged/server/server.go | 35 |
6 files changed, 392 insertions, 155 deletions
diff --git a/edge/pkg/edged/cadvisor/cadvisor_linux.go b/edge/pkg/edged/cadvisor/cadvisor_linux.go deleted file mode 100644 index 3915887af..000000000 --- a/edge/pkg/edged/cadvisor/cadvisor_linux.go +++ /dev/null @@ -1,98 +0,0 @@ -// +build cgo,linux - -/* -Copyright 2015 The Kubernetes Authors. -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. - -@CHANGELOG -KubeEdge Authors: To init a cadvisor client, implement the cadvisor interface by empty functions. -This file is derived from k8s kubelet code with pruned structures and interfaces -and changed most of the realization. -changes done are -1.For cadvisor.Interface is been implemented here. -*/ - -package cadvisor - -import ( - "github.com/google/cadvisor/events" - cadvisorapi "github.com/google/cadvisor/info/v1" - cadvisorapi2 "github.com/google/cadvisor/info/v2" - "k8s.io/kubernetes/pkg/kubelet/cadvisor" -) - -type cadvisorClient struct { - rootPath string -} - -// New creates new cadvisor client -func New(rootPath string) (cadvisor.Interface, error) { - return &cadvisorClient{ - rootPath: rootPath, - }, nil -} - -func (cadvisorClient) Start() error { - return nil -} - -func (cadvisorClient) DockerContainer(name string, req *cadvisorapi.ContainerInfoRequest) (cadvisorapi.ContainerInfo, error) { - return cadvisorapi.ContainerInfo{}, nil -} - -func (cadvisorClient) ContainerInfo(name string, req *cadvisorapi.ContainerInfoRequest) (*cadvisorapi.ContainerInfo, error) { - return nil, nil -} - -func (cadvisorClient) ContainerInfoV2(name string, options cadvisorapi2.RequestOptions) (map[string]cadvisorapi2.ContainerInfo, error) { - return nil, nil -} - -func (cadvisorClient) SubcontainerInfo(name string, req *cadvisorapi.ContainerInfoRequest) (map[string]*cadvisorapi.ContainerInfo, error) { - return nil, nil -} - -//MachineInfo implement by hard code, just for initing a cadvisor client, edged will not use this function to get machine info. -func (cadvisorClient) MachineInfo() (*cadvisorapi.MachineInfo, error) { - return &cadvisorapi.MachineInfo{ - NumCores: 4, - CpuFrequency: 3, - MemoryCapacity: 16000000, - HugePages: []cadvisorapi.HugePagesInfo{{PageSize: 4096, NumPages: 1024}}, - Filesystems: []cadvisorapi.FsInfo{}, - DiskMap: make(map[string]cadvisorapi.DiskInfo), - NetworkDevices: []cadvisorapi.NetInfo{}, - Topology: []cadvisorapi.Node{}, - CloudProvider: cadvisorapi.UnknownProvider, - InstanceType: cadvisorapi.UnknownInstance, - InstanceID: cadvisorapi.UnNamedInstance, - }, nil -} - -func (cadvisorClient) VersionInfo() (*cadvisorapi.VersionInfo, error) { - return nil, nil -} - -func (cadvisorClient) ImagesFsInfo() (cadvisorapi2.FsInfo, error) { - return cadvisorapi2.FsInfo{}, nil -} - -func (cadvisorClient) RootFsInfo() (cadvisorapi2.FsInfo, error) { - return cadvisorapi2.FsInfo{}, nil -} - -func (cadvisorClient) WatchEvents(request *events.Request) (*events.EventChannel, error) { - return nil, nil -} - -func (cadvisorClient) GetDirFsInfo(path string) (cadvisorapi2.FsInfo, error) { - return cadvisorapi2.FsInfo{}, nil -} diff --git a/edge/pkg/edged/containers/helpers.go b/edge/pkg/edged/containers/helpers.go index 7301a8064..374633f5e 100644 --- a/edge/pkg/edged/containers/helpers.go +++ b/edge/pkg/edged/containers/helpers.go @@ -23,12 +23,6 @@ Changes done are package containers -import ( - "time" - - kubecontainer "k8s.io/kubernetes/pkg/kubelet/container" -) - //KubeSourcesReady is blank structure just for function referencing type KubeSourcesReady struct{} @@ -36,16 +30,3 @@ type KubeSourcesReady struct{} func (s *KubeSourcesReady) AllReady() bool { return true } - -type containerRunner struct { -} - -func (c *containerRunner) RunInContainer(id kubecontainer.ContainerID, cmd []string, timeout time.Duration) ([]byte, error) { - return nil, nil -} - -//NewContainerRunner returns container manager object -// TODO: we didn't realized Run In container interface yet -func NewContainerRunner() kubecontainer.ContainerCommandRunner { - return &containerRunner{} -} diff --git a/edge/pkg/edged/edged.go b/edge/pkg/edged/edged.go index 0c77ddc5b..46c7079ec 100644 --- a/edge/pkg/edged/edged.go +++ b/edge/pkg/edged/edged.go @@ -52,6 +52,7 @@ import ( "k8s.io/klog" pluginwatcherapi "k8s.io/kubelet/pkg/apis/pluginregistration/v1" kubeletinternalconfig "k8s.io/kubernetes/pkg/kubelet/apis/config" + "k8s.io/kubernetes/pkg/kubelet/cadvisor" "k8s.io/kubernetes/pkg/kubelet/cm" "k8s.io/kubernetes/pkg/kubelet/cm/cpumanager" klconfigmap "k8s.io/kubernetes/pkg/kubelet/configmap" @@ -68,7 +69,9 @@ import ( "k8s.io/kubernetes/pkg/kubelet/prober" proberesults "k8s.io/kubernetes/pkg/kubelet/prober/results" "k8s.io/kubernetes/pkg/kubelet/remote" + serverstats "k8s.io/kubernetes/pkg/kubelet/server/stats" "k8s.io/kubernetes/pkg/kubelet/server/streaming" + "k8s.io/kubernetes/pkg/kubelet/stats" kubestatus "k8s.io/kubernetes/pkg/kubelet/status" "k8s.io/kubernetes/pkg/kubelet/util/format" "k8s.io/kubernetes/pkg/kubelet/util/queue" @@ -92,7 +95,6 @@ import ( "github.com/kubeedge/kubeedge/common/constants" "github.com/kubeedge/kubeedge/edge/pkg/common/modules" "github.com/kubeedge/kubeedge/edge/pkg/edged/apis" - "github.com/kubeedge/kubeedge/edge/pkg/edged/cadvisor" "github.com/kubeedge/kubeedge/edge/pkg/edged/clcm" edgedconfig "github.com/kubeedge/kubeedge/edge/pkg/edged/config" "github.com/kubeedge/kubeedge/edge/pkg/edged/containers" @@ -164,13 +166,14 @@ type podReady struct { podReadyLock sync.RWMutex } -//Define edged +// edged is the main edged implementation. type edged struct { - //dns config + // dns config dnsConfigurer *kubedns.Configurer hostname string namespace string nodeName string + runtimeCache kubecontainer.RuntimeCache interfaceName string uid types.UID nodeStatusUpdateFrequency time.Duration @@ -182,6 +185,7 @@ type edged struct { containerRuntime kubecontainer.Runtime podCache kubecontainer.Cache os kubecontainer.OSInterface + resourceAnalyzer serverstats.ResourceAnalyzer runtimeService internalapi.RuntimeService podManager podmanager.Manager pleg pleg.PodLifecycleEventGenerator @@ -212,13 +216,19 @@ type edged struct { configMapStore cache.Store workQueue queue.WorkQueue clcm clcm.ContainerLifecycleManager - //edged cgroup driver for container runtime + // edged cgroup driver for container runtime cgroupDriver string - //clusterDns dns + // clusterDns dns clusterDNS []net.IP // edge node IP nodeIP net.IP + // StatsProvider provides the node and the container stats. + *stats.StatsProvider + + // cAdvisor used for container information. + cadvisor cadvisor.Interface + // pluginmanager runs a set of asynchronous loops that figure out which // plugins need to be registered/unregistered based on this node and makes it so. pluginManager pluginmanager.PluginManager @@ -227,6 +237,13 @@ type edged struct { enable bool configMapManager klconfigmap.Manager + + // Cached MachineInfo returned by cadvisor. + machineInfo *cadvisorapi.MachineInfo + + dockerLegacyService dockershim.DockerLegacyService + // Optional, defaults to simple Docker implementation + runner kubecontainer.ContainerCommandRunner } // Register register edged @@ -257,7 +274,6 @@ func (e *edged) Enable() bool { func (e *edged) Start() { e.volumePluginMgr = NewInitializedVolumePluginMgr(e, ProbeVolumePlugins("")) - e.statusManager = status.NewManager(e.kubeClient, e.podManager, utilpod.NewPodDeleteSafety(), e.metaClient) if err := e.initializeModules(); err != nil { klog.Errorf("initialize module error: %v", err) os.Exit(1) @@ -285,7 +301,7 @@ func (e *edged) Start() { go e.volumeManager.Run(edgedutil.NewSourcesReady(), utilwait.NeverStop) go utilwait.Until(e.syncNodeStatus, e.nodeStatusUpdateFrequency, utilwait.NeverStop) - e.probeManager = prober.NewManager(e.statusManager, e.livenessManager, e.startupManager, containers.NewContainerRunner(), kubecontainer.NewRefManager(), record.NewEventRecorder()) + e.probeManager = prober.NewManager(e.statusManager, e.livenessManager, e.startupManager, e.runner, kubecontainer.NewRefManager(), record.NewEventRecorder()) e.pleg = pleg.NewGenericPLEG(e.containerRuntime, plegChannelCapacity, plegRelistPeriod, e.podCache, clock.RealClock{}) e.statusManager.Start() e.pleg.Start() @@ -297,7 +313,7 @@ func (e *edged) Start() { syncWorkQueueCh := time.NewTicker(syncWorkQueuePeriod) e.probeManager.Start() go e.syncLoopIteration(e.pleg.Watch(), housekeepingTicker.C, syncWorkQueueCh.C) - go e.server.ListenAndServe() + go e.server.ListenAndServe(e, e.resourceAnalyzer, true) e.imageGCManager.Start() e.StartGarbageCollection() @@ -343,6 +359,28 @@ func getRuntimeAndImageServices(remoteRuntimeEndpoint string, remoteImageEndpoin return rs, is, err } +func (e *edged) cgroupRoots() []string { + var cgroupRoots []string + + cgroupRoots = append(cgroupRoots, cm.NodeAllocatableRoot("", edgedconfig.Config.CGroupDriver)) + kubeletCgroup, err := cm.GetKubeletContainer("") + if err != nil { + klog.Warningf("failed to get the edged's cgroup: %v. Edged system container metrics may be missing.", err) + } else if kubeletCgroup != "" { + cgroupRoots = append(cgroupRoots, kubeletCgroup) + } + + runtimeCgroup, err := cm.GetRuntimeContainer(e.containerRuntimeName, "") + if err != nil { + klog.Warningf("failed to get the container runtime's cgroup: %v. Runtime system container metrics may be missing.", err) + } else if runtimeCgroup != "" { + // RuntimeCgroups is optional, so ignore if it isn't specified + cgroupRoots = append(cgroupRoots, runtimeCgroup) + } + + return cgroupRoots +} + //newEdged creates new edged object and initialises it func newEdged(enable bool) (*edged, error) { backoff := flowcontrol.NewBackOff(backOffPeriod, MaxContainerBackOff) @@ -362,6 +400,7 @@ func newEdged(enable bool) (*edged, error) { nodeName: edgedconfig.Config.HostnameOverride, interfaceName: edgedconfig.Config.InterfaceName, namespace: edgedconfig.Config.RegisterNodeNamespace, + containerRuntimeName: edgedconfig.Config.RuntimeType, gpuPluginEnabled: edgedconfig.Config.GPUPluginEnabled, cgroupDriver: edgedconfig.Config.CGroupDriver, concurrentConsumers: edgedconfig.Config.ConcurrentConsumers, @@ -453,7 +492,14 @@ func newEdged(enable bool) (*edged, error) { if err := server.Start(); err != nil { return nil, err } - + // Create dockerLegacyService when the logging driver is not supported. + supported, err := ds.IsCRISupportedLogDriver() + if err != nil { + return nil, err + } + if !supported { + ed.dockerLegacyService = ds + } } ed.clusterDNS = convertStrToIP(edgedconfig.Config.ClusterDNS) ed.dnsConfigurer = kubedns.NewConfigurer(recorder, @@ -480,15 +526,27 @@ func newEdged(enable bool) (*edged, error) { ed.clcm, err = clcm.NewContainerLifecycleManager(DefaultRootDir) - var machineInfo cadvisorapi.MachineInfo - machineInfo.MemoryCapacity = uint64(edgedconfig.Config.EdgedMemoryCapacity) + useLegacyCadvisorStats := cadvisor.UsingLegacyCadvisorStats(edgedconfig.Config.RuntimeType, edgedconfig.Config.RemoteRuntimeEndpoint) + imageFsInfoProvider := cadvisor.NewImageFsInfoProvider(edgedconfig.Config.RuntimeType, edgedconfig.Config.RemoteRuntimeEndpoint) + cadvisorInterface, err := cadvisor.New(imageFsInfoProvider, ed.rootDirectory, ed.cgroupRoots(), useLegacyCadvisorStats) + if err != nil { + return nil, err + } + ed.cadvisor = cadvisorInterface + + machineInfo, err := ed.cadvisor.MachineInfo() + if err != nil { + return nil, err + } + ed.machineInfo = machineInfo + containerRuntime, err := kuberuntime.NewKubeGenericRuntimeManager( recorder, ed.livenessManager, ed.startupManager, "", containerRefManager, - &machineInfo, + machineInfo, ed, ed.os, ed, @@ -509,7 +567,6 @@ func newEdged(enable bool) (*edged, error) { return nil, fmt.Errorf("New generic runtime manager failed, err: %s", err.Error()) } - cadvisorInterface, err := cadvisor.New("") containerManager, err := cm.NewContainerManager(mount.New(""), cadvisorInterface, cm.NodeConfig{ @@ -517,6 +574,7 @@ func newEdged(enable bool) (*edged, error) { SystemCgroupsName: edgedconfig.Config.CGroupDriver, KubeletCgroupsName: edgedconfig.Config.CGroupDriver, ContainerRuntime: edgedconfig.Config.RuntimeType, + CgroupsPerQOS: false, KubeletRootDir: DefaultRootDir, ExperimentalCPUManagerPolicy: string(cpumanager.PolicyNone), CgroupsPerQOS: edgedconfig.Config.CgroupsPerQOS, @@ -528,10 +586,42 @@ func newEdged(enable bool) (*edged, error) { if err != nil { return nil, fmt.Errorf("init container manager failed with error: %v", err) } + ed.containerRuntime = containerRuntime - ed.containerRuntimeName = edgedconfig.RemoteContainerRuntime + ed.runner = containerRuntime ed.containerManager = containerManager ed.runtimeService = runtimeService + + runtimeCache, err := kubecontainer.NewRuntimeCache(ed.containerRuntime) + if err != nil { + return nil, err + } + ed.runtimeCache = runtimeCache + + ed.resourceAnalyzer = serverstats.NewResourceAnalyzer(ed, time.Minute) + + ed.statusManager = status.NewManager(ed.kubeClient, ed.podManager, utilpod.NewPodDeleteSafety(), ed.metaClient) + + if useLegacyCadvisorStats { + ed.StatsProvider = stats.NewCadvisorStatsProvider( + cadvisorInterface, + ed.resourceAnalyzer, + ed.podManager, + ed.runtimeCache, + ed.containerRuntime, + ed.statusManager) + } else { + ed.StatsProvider = stats.NewCRIStatsProvider( + cadvisorInterface, + ed.resourceAnalyzer, + ed.podManager, + ed.runtimeCache, + ed.runtimeService, + imageService, + stats.NewLogMetricsService(), + kubecontainer.RealOS{}) + } + imageGCManager, err := images.NewImageGCManager( ed.containerRuntime, statsProvider, @@ -559,17 +649,33 @@ func newEdged(enable bool) (*edged, error) { } func (e *edged) initializeModules() error { + if err := e.cadvisor.Start(); err != nil { + // Fail kubelet and rely on the babysitter to retry starting kubelet. + // TODO(random-liu): Add backoff logic in the babysitter + klog.Fatalf("Failed to start cAdvisor %v", err) + } + + // trigger on-demand stats collection once so that we have capacity information for ephemeral storage. + // ignore any errors, since if stats collection is not successful, the container manager will fail to start below. + e.StatsProvider.GetCgroupStats("/", true) + + // Start container manager. node, err := e.initialNode() if err != nil { klog.Errorf("Failed to initialNode %v", err) return err } + // containerManager must start after cAdvisor because it needs filesystem capacity information err = e.containerManager.Start(node, e.GetActivePods, edgedutil.NewSourcesReady(), e.statusManager, e.runtimeService) if err != nil { - klog.Errorf("Failed to start device plugin manager %v", err) + klog.Errorf("Failed to start container manager, err: %v", err) return err } + + // Start resource analyzer + e.resourceAnalyzer.Start() + return nil } diff --git a/edge/pkg/edged/edged_getters.go b/edge/pkg/edged/edged_getters.go index 6984a3c99..f10100201 100644 --- a/edge/pkg/edged/edged_getters.go +++ b/edge/pkg/edged/edged_getters.go @@ -29,11 +29,17 @@ import ( "os" "path" "path/filepath" + "time" + + cadvisorapiv1 "github.com/google/cadvisor/info/v1" v1 "k8s.io/api/core/v1" "k8s.io/apimachinery/pkg/types" "k8s.io/klog" + "k8s.io/kubernetes/pkg/kubelet/cm" "k8s.io/kubernetes/pkg/kubelet/config" + "k8s.io/kubernetes/pkg/volume" + volumetypes "k8s.io/kubernetes/pkg/volume/util/types" "k8s.io/utils/mount" utilfile "k8s.io/utils/path" ) @@ -210,3 +216,87 @@ func (e *edged) podVolumeSubpathsDirExists(podUID types.UID) (bool, error) { func (e *edged) getPodVolumeSubpathsDir(podUID types.UID) string { return filepath.Join(e.getPodDir(podUID), config.DefaultKubeletVolumeSubpathsDirName) } + +// GetPodByName provides the first pod that matches namespace and name, as well +// as whether the pod was found. +func (e *edged) GetPodByName(namespace, name string) (*v1.Pod, bool) { + return e.podManager.GetPodByName(namespace, name) +} + +// GetNode returns the node info for the configured node name of this edged. +func (e *edged) GetNode() (*v1.Node, error) { + node := &v1.Node{} + node.Name = e.nodeName + return node, nil +} + +// GetNodeConfig returns the container manager node config. +func (e *edged) GetNodeConfig() cm.NodeConfig { + return e.containerManager.GetNodeConfig() +} + +func (e *edged) ListVolumesForPod(podUID types.UID) (map[string]volume.Volume, bool) { + volumesToReturn := make(map[string]volume.Volume) + podVolumes := e.volumeManager.GetMountedVolumesForPod( + volumetypes.UniquePodName(podUID)) + for outerVolumeSpecName, volume := range podVolumes { + // TODO: volume.Mounter could be nil if volume object is recovered + // from reconciler's sync state process. PR 33616 will fix this problem + // to create Mounter object when recovering volume state. + if volume.Mounter == nil { + continue + } + volumesToReturn[outerVolumeSpecName] = volume.Mounter + } + + return volumesToReturn, len(volumesToReturn) > 0 +} + +func (e *edged) GetPods() []*v1.Pod { + return e.podManager.GetPods() +} + +func (e *edged) GetPodCgroupRoot() string { + return e.containerManager.GetPodCgroupRoot() +} + +func (e *edged) GetPodByCgroupfs(cgroupfs string) (*v1.Pod, bool) { + pcm := e.containerManager.NewPodContainerManager() + if result, podUID := pcm.IsPodCgroup(cgroupfs); result { + return e.podManager.GetPodByUID(podUID) + } + return nil, false +} + +func (e *edged) GetVersionInfo() (*cadvisorapiv1.VersionInfo, error) { + return e.cadvisor.VersionInfo() +} + +func (e *edged) GetCachedMachineInfo() (*cadvisorapiv1.MachineInfo, error) { + return e.machineInfo, nil +} + +func (e *edged) GetRunningPods() ([]*v1.Pod, error) { + pods, err := e.runtimeCache.GetPods() + if err != nil { + return nil, err + } + + apiPods := make([]*v1.Pod, 0, len(pods)) + for _, pod := range pods { + apiPods = append(apiPods, pod.ToAPIPod()) + } + return apiPods, nil +} + +func (e *edged) ResyncInterval() time.Duration { + return time.Minute +} + +func (e *edged) GetHostname() string { + return e.hostname +} + +func (e *edged) LatestLoopEntryTime() time.Time { + return time.Time{} +} diff --git a/edge/pkg/edged/edged_pods.go b/edge/pkg/edged/edged_pods.go index 4094df095..19444da7a 100644 --- a/edge/pkg/edged/edged_pods.go +++ b/edge/pkg/edged/edged_pods.go @@ -29,9 +29,16 @@ package edged import ( "bytes" + "context" "fmt" + "io" "io/ioutil" + "k8s.io/kubernetes/pkg/kubelet/images" + "k8s.io/kubernetes/pkg/kubelet/server/portforward" + "k8s.io/kubernetes/pkg/kubelet/server/remotecommand" "net" + "net/http" + "net/url" "os" "path" "path/filepath" @@ -1144,3 +1151,165 @@ func (e *edged) getHostIPByInterface() (string, error) { } return "", fmt.Errorf("no ip and mask in this network card") } + +// findContainer finds and returns the container with the given pod ID, full name, and container name. +// It returns nil if not found. +func (e *edged) findContainer(podFullName string, podUID types.UID, containerName string) (*kubecontainer.Container, error) { + pods, err := e.containerRuntime.GetPods(false) + if err != nil { + return nil, err + } + // Resolve and type convert back again. + // We need the static pod UID but the kubecontainer API works with types.UID. + podUID = types.UID(e.podManager.TranslatePodUID(podUID)) + pod := kubecontainer.Pods(pods).FindPod(podFullName, podUID) + return pod.FindContainerByName(containerName), nil +} + +// RunInContainer runs a command in a container, returns the combined stdout, stderr as an array of bytes +func (e *edged) RunInContainer(podFullName string, podUID types.UID, containerName string, cmd []string) ([]byte, error) { + container, err := e.findContainer(podFullName, podUID, containerName) + if err != nil { + return nil, err + } + if container == nil { + return nil, fmt.Errorf("container not found (%q)", containerName) + } + // TODO(tallclair): Pass a proper timeout value. + return e.runner.RunInContainer(container.ID, cmd, 0) +} + +// GetKubeletContainerLogs returns logs from the container +// TODO: this method is returning logs of random container attempts, when it should be returning the most recent attempt +// or all of them. +func (e *edged) GetKubeletContainerLogs(ctx context.Context, podFullName, containerName string, logOptions *v1.PodLogOptions, stdout, stderr io.Writer) error { + // Pod workers periodically write status to statusManager. If status is not + // cached there, something is wrong (or kubelet just restarted and hasn't + // caught up yet). Just assume the pod is not ready yet. + name, namespace, err := kubecontainer.ParsePodFullName(podFullName) + if err != nil { + return fmt.Errorf("unable to parse pod full name %q: %v", podFullName, err) + } + + pod, ok := e.GetPodByName(namespace, name) + if !ok { + return fmt.Errorf("pod %q cannot be found - no logs available", name) + } + + podUID := pod.UID + if mirrorPod, ok := e.podManager.GetMirrorPodByPod(pod); ok { + podUID = mirrorPod.UID + } + podStatus, found := e.statusManager.GetPodStatus(podUID) + if !found { + // If there is no cached status, use the status from the + // apiserver. This is useful if kubelet has recently been + // restarted. + podStatus = pod.Status + } + + // TODO: Consolidate the logic here with kuberuntime.GetContainerLogs, here we convert container name to containerID, + // but inside kuberuntime we convert container id back to container name and restart count. + // TODO: After separate container log lifecycle management, we should get log based on the existing log files + // instead of container status. + containerID, err := e.validateContainerLogStatus(pod.Name, &podStatus, containerName, logOptions.Previous) + if err != nil { + return err + } + + // Do a zero-byte write to stdout before handing off to the container runtime. + // This ensures at least one Write call is made to the writer when copying starts, + // even if we then block waiting for log output from the container. + if _, err := stdout.Write([]byte{}); err != nil { + return err + } + + if e.dockerLegacyService != nil { + // dockerLegacyService should only be non-nil when we actually need it, so + // inject it into the runtimeService. + // TODO(random-liu): Remove this hack after deprecating unsupported log driver. + return e.dockerLegacyService.GetContainerLogs(ctx, pod, containerID, logOptions, stdout, stderr) + } + return e.containerRuntime.GetContainerLogs(ctx, pod, containerID, logOptions, stdout, stderr) +} + +func (e *edged) ServeLogs(w http.ResponseWriter, req *http.Request) { +} + +func (e *edged) GetExec(podFullName string, podUID types.UID, containerName string, cmd []string, streamOpts remotecommand.Options) (*url.URL, error) { + return nil, nil +} + +func (e *edged) GetAttach(podFullName string, podUID types.UID, containerName string, streamOpts remotecommand.Options) (*url.URL, error) { + return nil, nil +} + +func (e *edged) GetPortForward(podName, podNamespace string, podUID types.UID, portForwardOpts portforward.V4Options) (*url.URL, error) { + return nil, nil +} + +// validateContainerLogStatus returns the container ID for the desired container to retrieve logs for, based on the state +// of the container. The previous flag will only return the logs for the last terminated container, otherwise, the current +// running container is preferred over a previous termination. If info about the container is not available then a specific +// error is returned to the end user. +func (e *edged) validateContainerLogStatus(podName string, podStatus *v1.PodStatus, containerName string, previous bool) (containerID kubecontainer.ContainerID, err error) { + var cID string + + cStatus, found := podutil.GetContainerStatus(podStatus.ContainerStatuses, containerName) + if !found { + cStatus, found = podutil.GetContainerStatus(podStatus.InitContainerStatuses, containerName) + } + if !found && utilfeature.DefaultFeatureGate.Enabled(features.EphemeralContainers) { + cStatus, found = podutil.GetContainerStatus(podStatus.EphemeralContainerStatuses, containerName) + } + if !found { + return kubecontainer.ContainerID{}, fmt.Errorf("container %q in pod %q is not available", containerName, podName) + } + lastState := cStatus.LastTerminationState + waiting, running, terminated := cStatus.State.Waiting, cStatus.State.Running, cStatus.State.Terminated + + switch { + case previous: + if lastState.Terminated == nil || lastState.Terminated.ContainerID == "" { + return kubecontainer.ContainerID{}, fmt.Errorf("previous terminated container %q in pod %q not found", containerName, podName) + } + cID = lastState.Terminated.ContainerID + + case running != nil: + cID = cStatus.ContainerID + + case terminated != nil: + // in cases where the next container didn't start, terminated.ContainerID will be empty, so get logs from the lastState.Terminated. + if terminated.ContainerID == "" { + if lastState.Terminated != nil && lastState.Terminated.ContainerID != "" { + cID = lastState.Terminated.ContainerID + } else { + return kubecontainer.ContainerID{}, fmt.Errorf("container %q in pod %q is terminated", containerName, podName) + } + } else { + cID = terminated.ContainerID + } + + case lastState.Terminated != nil: + if lastState.Terminated.ContainerID == "" { + return kubecontainer.ContainerID{}, fmt.Errorf("container %q in pod %q is terminated", containerName, podName) + } + cID = lastState.Terminated.ContainerID + + case waiting != nil: + // output some info for the most common pending failures + switch reason := waiting.Reason; reason { + case images.ErrImagePull.Error(): + return kubecontainer.ContainerID{}, fmt.Errorf("container %q in pod %q is waiting to start: image can't be pulled", containerName, podName) + case images.ErrImagePullBackOff.Error(): + return kubecontainer.ContainerID{}, fmt.Errorf("container %q in pod %q is waiting to start: trying and failing to pull image", containerName, podName) + default: + return kubecontainer.ContainerID{}, fmt.Errorf("container %q in pod %q is waiting to start: %v", containerName, podName, reason) + } + default: + // unrecognized state + return kubecontainer.ContainerID{}, fmt.Errorf("container %q in pod %q is waiting to start - no logs yet", containerName, podName) + } + + return kubecontainer.ParseContainerID(cID), nil +} diff --git a/edge/pkg/edged/server/server.go b/edge/pkg/edged/server/server.go index 01c3ff3b9..f982d0dc4 100644 --- a/edge/pkg/edged/server/server.go +++ b/edge/pkg/edged/server/server.go @@ -1,14 +1,12 @@ package server import ( - "bytes" - "encoding/json" "net" "net/http" - "strconv" - "k8s.io/api/core/v1" "k8s.io/klog" + "k8s.io/kubernetes/pkg/kubelet/server" + "k8s.io/kubernetes/pkg/kubelet/server/stats" "github.com/kubeedge/kubeedge/edge/pkg/edged/podmanager" ) @@ -16,7 +14,7 @@ import ( //constants to define server address const ( ServerAddr = "127.0.0.1" - ServerPort = 10255 + ServerPort = "10255" ) //Server is object to define server @@ -31,24 +29,15 @@ func NewServer(podManager podmanager.Manager) *Server { } } -func (s *Server) getPodsHandler(w http.ResponseWriter, r *http.Request) { - var podList v1.PodList - pods := s.podManager.GetPods() - for _, pod := range pods { - podList.Items = append(podList.Items, *pod) - } - rspBodyBytes := new(bytes.Buffer) - json.NewEncoder(rspBodyBytes).Encode(podList) - w.Write(rspBodyBytes.Bytes()) -} - // ListenAndServe starts a HTTP server and sets up a listener on the given host/port -func (s *Server) ListenAndServe() { - klog.Infof("starting to listen on %s:%d", ServerAddr, ServerPort) - mux := http.NewServeMux() - mux.HandleFunc("/pods", s.getPodsHandler) - err := http.ListenAndServe(net.JoinHostPort(ServerAddr, strconv.FormatUint(uint64(ServerPort), 10)), mux) - if err != nil { - klog.Fatalf("run server: %v", err) +func (s *Server) ListenAndServe(host server.HostInterface, resourceAnalyzer stats.ResourceAnalyzer, enableCAdvisorJSONEndpoints bool) { + klog.Infof("starting to listen read-only on %s:%s", ServerAddr, ServerPort) + handler := server.NewServer(host, resourceAnalyzer, nil, enableCAdvisorJSONEndpoints, false, false, false, nil) + + server := &http.Server{ + Addr: net.JoinHostPort(ServerAddr, ServerPort), + Handler: &handler, + MaxHeaderBytes: 1 << 20, } + klog.Fatal(server.ListenAndServe()) } |
