summaryrefslogtreecommitdiff
path: root/edge
diff options
context:
space:
mode:
authorfisherxu <xufei40@huawei.com>2020-03-28 18:41:25 +0800
committerfisherxu <xufei40@huawei.com>2020-04-17 17:50:46 +0800
commit062b3ddd6b70bad35057359a273d443c5d972d72 (patch)
treeb8ab42d2eec8a736349a2fc7460176b3906f7fce /edge
parentadd vendor (diff)
downloadkubeedge-062b3ddd6b70bad35057359a273d443c5d972d72.tar.gz
add cadvisor
Diffstat (limited to 'edge')
-rw-r--r--edge/pkg/edged/cadvisor/cadvisor_linux.go98
-rw-r--r--edge/pkg/edged/containers/helpers.go19
-rw-r--r--edge/pkg/edged/edged.go136
-rw-r--r--edge/pkg/edged/edged_getters.go90
-rw-r--r--edge/pkg/edged/edged_pods.go169
-rw-r--r--edge/pkg/edged/server/server.go35
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())
}