feat(host): container metrics (#20257)

This commit is contained in:
Zexi Li
2024-05-16 16:00:29 +08:00
committed by GitHub
parent 52eb75ecf4
commit b739cbd0dc
825 changed files with 145711 additions and 255 deletions
-4
View File
@@ -459,10 +459,6 @@ func (s *SKVMGuestInstance) IsStopping() bool {
return s.stopping
}
func (s *SKVMGuestInstance) IsValid() bool {
return s.Desc != nil && s.Desc.Uuid != ""
}
func (s *SKVMGuestInstance) getStateFilePathRootPrefix() string {
return path.Join(s.HomeDir(), STATE_FILE_PREFIX)
}
+10
View File
@@ -35,9 +35,11 @@ import (
)
type GuestRuntimeInstance interface {
GetHypervisor() string
GetName() string
GetInitialId() string
GetId() string
IsValid() bool
HomeDir() string
GetDesc() *desc.SGuestDesc
SetDesc(guestDesc *desc.SGuestDesc)
@@ -84,6 +86,14 @@ func newBaseGuestInstance(id string, manager *SGuestManager, hypervisor string)
}
}
func (s *sBaseGuestInstance) GetHypervisor() string {
return s.Hypervisor
}
func (s *sBaseGuestInstance) IsValid() bool {
return s.Desc != nil && s.Desc.Uuid != ""
}
func (s *sBaseGuestInstance) GetInitialId() string {
return s.Id
}
+1 -1
View File
@@ -112,7 +112,7 @@ func (host *SHostService) RunService() {
log.Fatalf("Guest manager Bootstrap %s", err)
}
// hostmetrics after guestmanager bootstrap
hostmetrics.Init()
hostmetrics.Init(hostInstance.GetContainerStatsProvider())
hostmetrics.Start()
host.initHandlers(app)
+24 -2
View File
@@ -71,6 +71,8 @@ import (
"yunion.io/x/onecloud/pkg/util/netutils2"
"yunion.io/x/onecloud/pkg/util/ovnutils"
"yunion.io/x/onecloud/pkg/util/pod"
"yunion.io/x/onecloud/pkg/util/pod/cadvisor"
"yunion.io/x/onecloud/pkg/util/pod/stats"
"yunion.io/x/onecloud/pkg/util/procutils"
"yunion.io/x/onecloud/pkg/util/qemutils"
"yunion.io/x/onecloud/pkg/util/sysutils"
@@ -118,8 +120,9 @@ type SHostInfo struct {
IoScheduler string
cri pod.CRI
containerCPUMap *pod.HostContainerCPUMap
cri pod.CRI
containerCPUMap *pod.HostContainerCPUMap
containerStatsProvier stats.ContainerStatsProvider
}
func (h *SHostInfo) GetContainerDeviceConfigurationFilePath() string {
@@ -229,6 +232,9 @@ func (h *SHostInfo) Init() error {
if err := h.initContainerCPUMap(h.sysinfo.Topology); err != nil {
return errors.Wrap(err, "init container cpu map")
}
if err := h.startContainerStatsProvider(h.cri); err != nil {
return errors.Wrap(err, "start container stats provider")
}
}
return nil
@@ -258,6 +264,18 @@ func (h *SHostInfo) initContainerCPUMap(topo *hostapi.HostTopology) error {
return nil
}
func (h *SHostInfo) startContainerStatsProvider(cri pod.CRI) error {
ca, err := cadvisor.New(nil, "/opt/cloud/workspace", []string{"cloudpods"})
if err != nil {
return errors.Wrap(err, "new cadvisor")
}
if err := ca.Start(); err != nil {
return errors.Wrap(err, "start cadvisor")
}
h.containerStatsProvier = stats.NewCRIContainerStatsProvider(ca, cri.GetRuntimeClient(), cri.GetImageClient())
return nil
}
func (h *SHostInfo) GetCRI() pod.CRI {
return h.cri
}
@@ -266,6 +284,10 @@ func (h *SHostInfo) GetContainerCPUMap() *pod.HostContainerCPUMap {
return h.containerCPUMap
}
func (h *SHostInfo) GetContainerStatsProvider() stats.ContainerStatsProvider {
return h.containerStatsProvier
}
func (h *SHostInfo) setupOvnChassis() error {
opts := &options.HostOptions
if opts.BridgeDriver != hostbridge.DRV_OPEN_VSWITCH {
@@ -0,0 +1,241 @@
package hostmetrics
import (
"fmt"
"strings"
"time"
"yunion.io/x/onecloud/pkg/util/pod/stats"
)
const (
// Cumulative cpu time consumed by the container in core-seconds
CPU_USAGE_SECONDS_TOTAL = "usage_seconds_total"
// cpu usage rate
CPU_USAGE_RATE = "usage_rate"
// current working set of memory size in bytes
MEMORY_WORKING_SET_BYTES = "working_set_bytes"
// memory usage rate
MEMORY_USAGE_RATE = "usage_rate"
)
type PodMetrics struct {
PodCpu *PodCpuMetric `json:"pod_cpu"`
PodMemory *PodMemoryMetric `json:"pod_memory"`
Containers []*ContainerMetrics `json:"containers"`
}
type PodMetricMeta struct {
Time time.Time
}
func NewPodMetricMeta(time time.Time) PodMetricMeta {
return PodMetricMeta{Time: time}
}
func (m PodMetricMeta) GetTag() map[string]string {
return nil
}
type PodCpuMetric struct {
PodMetricMeta
CpuUsageSecondsTotal float64 `json:"cpu_usage_seconds_total"`
CpuUsageRate *float64 `json:"cpu_usage_rate"`
}
func (m PodCpuMetric) GetName() string {
return "pod_cpu"
}
func (m PodCpuMetric) ToMap() map[string]interface{} {
ret := map[string]interface{}{
CPU_USAGE_SECONDS_TOTAL: m.CpuUsageSecondsTotal,
}
if m.CpuUsageRate != nil {
ret[CPU_USAGE_RATE] = *m.CpuUsageRate
}
return ret
}
type PodMemoryMetric struct {
PodMetricMeta
MemoryWorkingSetBytes float64 `json:"memory_working_set_bytes"`
MemoryUsageRate float64 `json:"memory_usage_rate"`
}
func (m PodMemoryMetric) GetName() string {
return "pod_mem"
}
func (m PodMemoryMetric) ToMap() map[string]interface{} {
return map[string]interface{}{
MEMORY_WORKING_SET_BYTES: m.MemoryWorkingSetBytes,
MEMORY_USAGE_RATE: m.MemoryUsageRate,
}
}
type ContainerMetrics struct {
ContainerCpu *ContainerCpuMetric `json:"container_cpu"`
ContainerMemory *ContainerMemoryMetric `json:"container_memory"`
}
type ContainerMetricMeta struct {
PodMetricMeta
ContainerId string `json:"container_id"`
ContainerName string `json:"container_name"`
PodId string `json:"pod_id"`
}
func (m ContainerMetricMeta) GetTag() map[string]string {
ret := map[string]string{
"pod_id": m.PodId,
"container_name": m.ContainerName,
}
if m.ContainerId != "" {
ret["container_id"] = m.ContainerId
}
return ret
}
type ContainerMemoryMetric struct {
ContainerMetricMeta
MemoryWorkingSetBytes float64 `json:"memory_working_set_bytes"`
MemoryUsageRate float64 `json:"memory_usage_rate"`
}
func (m ContainerMemoryMetric) GetName() string {
return "container_mem"
}
func (m *ContainerMemoryMetric) ToMap() map[string]interface{} {
return map[string]interface{}{
MEMORY_WORKING_SET_BYTES: m.MemoryWorkingSetBytes,
MEMORY_USAGE_RATE: m.MemoryUsageRate,
}
}
type ContainerCpuMetric struct {
ContainerMetricMeta
CpuUsageSecondsTotal float64 `json:"cpu_usage_seconds_total"`
CpuUsageRate *float64 `json:"cpu_usage_rate"`
}
func (m ContainerCpuMetric) GetName() string {
return "container_cpu"
}
func (m *ContainerCpuMetric) ToMap() map[string]interface{} {
ret := map[string]interface{}{
CPU_USAGE_SECONDS_TOTAL: m.CpuUsageSecondsTotal,
}
if m.CpuUsageRate != nil {
ret[CPU_USAGE_RATE] = *m.CpuUsageRate
}
return ret
}
func GetPodStatsById(stats []stats.PodStats, podId string) *stats.PodStats {
for _, stat := range stats {
if stat.PodRef.UID == podId {
tmp := stat
return &tmp
}
}
return nil
}
func (s *SGuestMonitorCollector) collectPodMetrics(gm *SGuestMonitor, prevUsage *GuestMetrics) *GuestMetrics {
gmData := new(GuestMetrics)
gmData.PodMetrics = gm.PodMetrics(prevUsage)
return gmData
}
func NewContainerMetricMeta(serverId string, containerId string, containerName string, time time.Time) ContainerMetricMeta {
return ContainerMetricMeta{
PodMetricMeta: NewPodMetricMeta(time),
ContainerId: containerId,
PodId: serverId,
ContainerName: containerName,
}
}
func (m *SGuestMonitor) HasPodMetrics() bool {
return m.podStat != nil
}
func (m *SGuestMonitor) PodMetrics(prevUsage *GuestMetrics) *PodMetrics {
stat := m.podStat
podCpu := &PodCpuMetric{
PodMetricMeta: NewPodMetricMeta(stat.CPU.Time.Time),
CpuUsageSecondsTotal: float64(*stat.CPU.UsageCoreNanoSeconds) / float64(time.Second),
}
hasPrevUsage := prevUsage != nil && prevUsage.PodMetrics != nil
if hasPrevUsage {
pmPodCpu := prevUsage.PodMetrics.PodCpu
val := (podCpu.CpuUsageSecondsTotal - pmPodCpu.CpuUsageSecondsTotal) / podCpu.Time.Sub(pmPodCpu.Time).Seconds() * 100
podCpu.CpuUsageRate = &val
}
podMemory := &PodMemoryMetric{
MemoryWorkingSetBytes: float64(*stat.Memory.WorkingSetBytes),
MemoryUsageRate: (float64(*stat.Memory.WorkingSetBytes) / float64(m.MemMB*1024*1024)) * 100,
}
containers := make([]*ContainerMetrics, 0)
for _, ctr := range stat.Containers {
meta := NewContainerMetricMeta(m.Id, "", ctr.Name, ctr.CPU.Time.Time)
cm := &ContainerMetrics{
ContainerCpu: &ContainerCpuMetric{
ContainerMetricMeta: meta,
CpuUsageSecondsTotal: float64(*ctr.CPU.UsageCoreNanoSeconds) / float64(time.Second),
},
ContainerMemory: &ContainerMemoryMetric{
ContainerMetricMeta: meta,
MemoryWorkingSetBytes: float64(*ctr.Memory.WorkingSetBytes),
MemoryUsageRate: (float64(*ctr.Memory.WorkingSetBytes) / float64(m.MemMB*1024*1024)) * 100,
},
}
if hasPrevUsage {
for _, prevCtr := range prevUsage.PodMetrics.Containers {
if prevCtr.ContainerCpu.ContainerName == ctr.Name {
val := (cm.ContainerCpu.CpuUsageSecondsTotal - prevCtr.ContainerCpu.CpuUsageSecondsTotal) / cm.ContainerCpu.Time.Sub(prevCtr.ContainerCpu.Time).Seconds() * 100
cm.ContainerCpu.CpuUsageRate = &val
break
}
}
}
containers = append(containers, cm)
}
return &PodMetrics{
PodCpu: podCpu,
PodMemory: podMemory,
Containers: containers,
}
}
type iPodMetric interface {
GetName() string
GetTag() map[string]string
ToMap() map[string]interface{}
}
func (d *GuestMetrics) toPodTelegrafData(tagStr string) []string {
m := d.PodMetrics
ims := []iPodMetric{m.PodCpu, m.PodMemory}
for _, c := range m.Containers {
ims = append(ims, c.ContainerCpu)
ims = append(ims, c.ContainerMemory)
}
res := []string{}
for _, im := range ims {
tagMap := im.GetTag()
if len(tagMap) != 0 {
var newTagArr []string
for k, v := range tagMap {
newTagArr = append(newTagArr, fmt.Sprintf("%s=%s", k, v))
}
tagStr = strings.Join([]string{tagStr, strings.Join(newTagArr, ",")}, ",")
}
res = append(res, fmt.Sprintf("%s,%s %s", im.GetName(), tagStr, d.mapToStatStr(im.ToMap())))
}
return res
}
+155 -63
View File
@@ -37,6 +37,7 @@ import (
"yunion.io/x/onecloud/pkg/hostman/hostinfo/hostconsts"
"yunion.io/x/onecloud/pkg/hostman/options"
"yunion.io/x/onecloud/pkg/util/fileutils2"
"yunion.io/x/onecloud/pkg/util/pod/stats"
)
const (
@@ -53,9 +54,9 @@ type SHostMetricsCollector struct {
var hostMetricsCollector *SHostMetricsCollector
func Init() {
func Init(csp stats.ContainerStatsProvider) {
if hostMetricsCollector == nil {
hostMetricsCollector = NewHostMetricsCollector()
hostMetricsCollector = NewHostMetricsCollector(csp)
}
}
@@ -140,25 +141,27 @@ func (m *SHostMetricsCollector) collectReportData() string {
return m.guestMonitor.CollectReportData()
}
func NewHostMetricsCollector() *SHostMetricsCollector {
func NewHostMetricsCollector(csp stats.ContainerStatsProvider) *SHostMetricsCollector {
return &SHostMetricsCollector{
ReportInterval: options.HostOptions.ReportInterval,
waitingReportData: make([]string, 0),
guestMonitor: NewGuestMonitorCollector(),
guestMonitor: NewGuestMonitorCollector(csp),
}
}
type SGuestMonitorCollector struct {
monitors map[string]*SGuestMonitor
prevPids map[string]int
prevReportData map[string]*GuestMetrics
monitors map[string]*SGuestMonitor
prevPids map[string]int
prevReportData map[string]*GuestMetrics
containerStatsProvider stats.ContainerStatsProvider
}
func NewGuestMonitorCollector() *SGuestMonitorCollector {
func NewGuestMonitorCollector(csp stats.ContainerStatsProvider) *SGuestMonitorCollector {
return &SGuestMonitorCollector{
monitors: make(map[string]*SGuestMonitor, 0),
prevPids: make(map[string]int, 0),
prevReportData: make(map[string]*GuestMetrics, 0),
monitors: make(map[string]*SGuestMonitor, 0),
prevPids: make(map[string]int, 0),
prevReportData: make(map[string]*GuestMetrics, 0),
containerStatsProvider: csp,
}
}
@@ -166,41 +169,72 @@ func (s *SGuestMonitorCollector) GetGuests() map[string]*SGuestMonitor {
var err error
gms := make(map[string]*SGuestMonitor, 0)
guestmanager := guestman.GetGuestManager()
var podStats []stats.PodStats = nil
guestmanager.Servers.Range(func(k, v interface{}) bool {
guest, ok := v.(*guestman.SKVMGuestInstance)
instance, ok := v.(guestman.GuestRuntimeInstance)
if !ok {
return false
}
if !guest.IsValid() {
if !instance.IsValid() {
return false
}
pid := guest.GetPid()
if pid > 0 {
guestName := guest.Desc.Name
guestId := guest.GetId()
nicsDesc := guest.Desc.Nics
vcpuCount := guest.Desc.Cpu
gm, ok := s.monitors[guestId]
if ok && gm.Pid == pid {
delete(s.monitors, guestId)
gm.UpdateVmName(guestName)
gm.UpdateNicsDesc(nicsDesc)
gm.UpdateCpuCount(int(vcpuCount))
} else {
delete(s.monitors, guestId)
gm, err = NewGuestMonitor(guestName, guestId, pid, nicsDesc, int(vcpuCount))
hypervisor := instance.GetHypervisor()
guestId := instance.GetId()
guestName := instance.GetDesc().Name
nicsDesc := instance.GetDesc().Nics
vcpuCount := instance.GetDesc().Cpu
switch hypervisor {
case compute.HYPERVISOR_KVM:
guest := instance.(*guestman.SKVMGuestInstance)
pid := guest.GetPid()
if pid > 0 {
gm, ok := s.monitors[guestId]
if ok && gm.Pid == pid {
delete(s.monitors, guestId)
gm.UpdateVmName(guestName)
gm.UpdateNicsDesc(nicsDesc)
gm.UpdateCpuCount(int(vcpuCount))
gm.MemMB = instance.GetDesc().Mem
} else {
delete(s.monitors, guestId)
gm, err = NewGuestMonitor(guestName, guestId, pid, nicsDesc, int(vcpuCount))
if err != nil {
log.Errorf("NewGuestMonitor for %s(%s), pid: %d, nics: %#v", guestName, guestId, pid, nicsDesc)
return true
}
}
gm.ScalingGroupId = guest.GetDesc().ScalingGroupId
gm.Tenant = guest.GetDesc().Tenant
gm.TenantId = guest.GetDesc().TenantId
gm.DomainId = guest.GetDesc().DomainId
gm.ProjectDomain = guest.GetDesc().ProjectDomain
gms[guestId] = gm
}
return true
case compute.HYPERVISOR_POD:
if podStats == nil {
var err error
podStats, err = s.containerStatsProvider.ListPodCPUAndMemoryStats()
if err != nil {
log.Errorf("NewGuestMonitor for %s(%s), pid: %d, nics: %#v", guestName, guestId, pid, nicsDesc)
log.Errorf("ListPodCPUAndMemoryStats: %s", err)
return true
}
}
gm.ScalingGroupId = guest.Desc.ScalingGroupId
gm.Tenant = guest.Desc.Tenant
gm.TenantId = guest.Desc.TenantId
gm.DomainId = guest.Desc.DomainId
gm.ProjectDomain = guest.Desc.ProjectDomain
gms[guestId] = gm
podStat := GetPodStatsById(podStats, guestId)
if podStat != nil {
gm, err := NewGuestPodMonitor(guestName, guestId, podStat, nicsDesc, int(vcpuCount))
if err != nil {
return true
}
gm.UpdateByInstance(instance)
gms[guestId] = gm
return true
} else {
delete(s.monitors, guestId)
}
}
return true
})
@@ -360,10 +394,30 @@ func (s *SGuestMonitorCollector) cleanedPrevData(gms map[string]*SGuestMonitor)
}
type GuestMetrics struct {
VmCpu *CpuMetric `json:"vm_cpu"`
VmMem *MemMetric `json:"vm_mem"`
VmNetio []*NetIOMetric `json:"vm_netio"`
VmDiskio *DiskIOMetric `json:"vm_diskio"`
VmCpu *CpuMetric `json:"vm_cpu"`
VmMem *MemMetric `json:"vm_mem"`
VmNetio []*NetIOMetric `json:"vm_netio"`
VmDiskio *DiskIOMetric `json:"vm_diskio"`
PodMetrics *PodMetrics `json:"pod_metrics"`
}
func (d *GuestMetrics) mapToStatStr(m map[string]interface{}) string {
var statArr = []string{}
for k, v := range m {
statArr = append(statArr, fmt.Sprintf("%s=%v", k, v))
}
return strings.Join(statArr, ",")
}
func (d *GuestMetrics) toVmTelegrafData(tagStr string) []string {
var res = []string{}
res = append(res, fmt.Sprintf("%s,%s %s", "vm_cpu", tagStr, d.mapToStatStr(d.VmCpu.ToMap())))
res = append(res, fmt.Sprintf("%s,%s %s", "vm_mem", tagStr, d.mapToStatStr(d.VmMem.ToMap())))
res = append(res, fmt.Sprintf("%s,%s %s", "vm_diskio", tagStr, d.mapToStatStr(d.VmDiskio.ToMap())))
for i := range d.VmNetio {
res = append(res, fmt.Sprintf("%s,%s %s", "vm_netio", tagStr, d.mapToStatStr(d.VmNetio[i].ToMap())))
}
return res
}
func (d *GuestMetrics) toTelegrafData(tags map[string]string) []string {
@@ -372,23 +426,11 @@ func (d *GuestMetrics) toTelegrafData(tags map[string]string) []string {
tagArr = append(tagArr, fmt.Sprintf("%s=%s", k, strings.ReplaceAll(v, " ", "+")))
}
tagStr := strings.Join(tagArr, ",")
mapToStatStr := func(m map[string]interface{}) string {
var statArr = []string{}
for k, v := range m {
statArr = append(statArr, fmt.Sprintf("%s=%v", k, v))
}
return strings.Join(statArr, ",")
if d.PodMetrics == nil {
return d.toVmTelegrafData(tagStr)
} else {
return d.toPodTelegrafData(tagStr)
}
var res = []string{}
res = append(res, fmt.Sprintf("%s,%s %s", "vm_cpu", tagStr, mapToStatStr(d.VmCpu.ToMap())))
res = append(res, fmt.Sprintf("%s,%s %s", "vm_mem", tagStr, mapToStatStr(d.VmMem.ToMap())))
res = append(res, fmt.Sprintf("%s,%s %s", "vm_diskio", tagStr, mapToStatStr(d.VmDiskio.ToMap())))
for i := range d.VmNetio {
res = append(res, fmt.Sprintf("%s,%s %s", "vm_netio", tagStr, mapToStatStr(d.VmNetio[i].ToMap())))
}
return res
}
func (s *SGuestMonitorCollector) collectGmReport(
@@ -397,6 +439,15 @@ func (s *SGuestMonitorCollector) collectGmReport(
if prevUsage == nil {
prevUsage = new(GuestMetrics)
}
if !gm.HasPodMetrics() {
return s.collectGuestMetrics(gm, prevUsage)
} else {
return s.collectPodMetrics(gm, prevUsage)
}
}
func (s *SGuestMonitorCollector) collectGuestMetrics(gm *SGuestMonitor, prevUsage *GuestMetrics) *GuestMetrics {
gmData := new(GuestMetrics)
gmData.VmCpu = gm.Cpu()
gmData.VmMem = gm.Mem()
@@ -480,6 +531,7 @@ type SGuestMonitor struct {
Pid int
Nics []*desc.SGuestNetwork
CpuCnt int
MemMB int64
Ip string
Process *process.Process
ScalingGroupId string
@@ -487,19 +539,59 @@ type SGuestMonitor struct {
TenantId string
DomainId string
ProjectDomain string
podStat *stats.PodStats
}
func NewGuestMonitor(name, id string, pid int, nics []*desc.SGuestNetwork, cpuCount int,
) (*SGuestMonitor, error) {
var ip string
if len(nics) >= 1 {
ip = nics[0].Ip
}
func NewGuestMonitor(name, id string, pid int, nics []*desc.SGuestNetwork, cpuCount int) (*SGuestMonitor, error) {
proc, err := process.NewProcess(int32(pid))
if err != nil {
return nil, err
}
return &SGuestMonitor{name, id, pid, nics, cpuCount, ip, proc, "", "", "", "", ""}, nil
return newGuestMonitor(name, id, proc, nics, cpuCount)
}
func NewGuestPodMonitor(name, id string, stat *stats.PodStats, nics []*desc.SGuestNetwork, cpuCount int) (*SGuestMonitor, error) {
m, err := newGuestMonitor(name, id, nil, nics, cpuCount)
if err != nil {
return nil, errors.Wrap(err, "new pod GuestMonitor")
}
m.podStat = stat
return m, nil
}
func newGuestMonitor(name, id string, proc *process.Process, nics []*desc.SGuestNetwork, cpuCount int) (*SGuestMonitor, error) {
var ip string
if len(nics) >= 1 {
ip = nics[0].Ip
}
pid := 0
if proc != nil {
pid = int(proc.Pid)
}
return &SGuestMonitor{
Name: name,
Id: id,
Pid: pid,
Nics: nics,
CpuCnt: cpuCount,
Ip: ip,
Process: proc,
}, nil
}
func (m *SGuestMonitor) UpdateByInstance(instance guestman.GuestRuntimeInstance) {
guestName := instance.GetDesc().Name
nicsDesc := instance.GetDesc().Nics
vcpuCount := instance.GetDesc().Cpu
m.UpdateVmName(guestName)
m.UpdateNicsDesc(nicsDesc)
m.UpdateCpuCount(int(vcpuCount))
m.MemMB = instance.GetDesc().Mem
m.ScalingGroupId = instance.GetDesc().ScalingGroupId
m.Tenant = instance.GetDesc().Tenant
m.TenantId = instance.GetDesc().TenantId
m.DomainId = instance.GetDesc().DomainId
m.ProjectDomain = instance.GetDesc().ProjectDomain
}
func (m *SGuestMonitor) SetNicDown(index int) {
+167
View File
@@ -0,0 +1,167 @@
package cadvisor
import (
"flag"
"fmt"
"net/http"
"os"
"path"
"time"
"github.com/google/cadvisor/cache/memory"
cadvisormetrics "github.com/google/cadvisor/container"
"github.com/google/cadvisor/container/containerd"
_ "github.com/google/cadvisor/container/containerd/install"
"github.com/google/cadvisor/events"
cadvisorapi "github.com/google/cadvisor/info/v1"
cadvisorapiv2 "github.com/google/cadvisor/info/v2"
"github.com/google/cadvisor/manager"
"github.com/google/cadvisor/utils/sysfs"
"yunion.io/x/log"
"yunion.io/x/pkg/errors"
)
const (
// The amount of time for which to keep stats in memory.
statsCacheDuration = 2 * time.Minute
maxHousekeepingInterval = 15 * time.Second
defaultHousekeepingInterval = 10 * time.Second
allowDynamicHousekeeping = true
)
func init() {
ep := "/var/run/onecloud/containerd/containerd.sock"
containerd.ArgContainerdEndpoint = &ep
// REF: k8s.io/kubernetes/pkg/kubelet/cadvisor/cadvisor_linux.go
// override cadvisor flag defaults.
flagOverrides := map[string]string{
// Override the default cadvisor housekeeping interval.
"housekeeping_interval": defaultHousekeepingInterval.String(),
// Disable event storage by default.
"event_storage_event_limit": "default=0",
"event_storage_age_limit": "default=0",
}
for name, defaultValue := range flagOverrides {
if f := flag.Lookup(name); f != nil {
f.DefValue = defaultValue
f.Value.Set(defaultValue)
} else {
log.Errorf("Expected cAdvisor flag %q not found", name)
}
}
}
type cadvisorClient struct {
manager.Manager
rootPath string
imageFsInfoProvider ImageFsInfoProvider
}
func New(imageFsInfoProvider ImageFsInfoProvider, rootPath string, cgroupRoots []string) (Interface, error) {
includedMetrics := cadvisormetrics.MetricSet{
cadvisormetrics.CpuUsageMetrics: struct{}{},
cadvisormetrics.MemoryUsageMetrics: struct{}{},
cadvisormetrics.CpuLoadMetrics: struct{}{},
cadvisormetrics.DiskIOMetrics: struct{}{},
cadvisormetrics.NetworkUsageMetrics: struct{}{},
cadvisormetrics.AcceleratorUsageMetrics: struct{}{},
cadvisormetrics.AppMetrics: struct{}{},
cadvisormetrics.ProcessMetrics: struct{}{},
cadvisormetrics.DiskUsageMetrics: struct{}{},
}
duration := maxHousekeepingInterval
allowDynamic := allowDynamicHousekeeping
housekeepingConfig := manager.HouskeepingConfig{
Interval: &duration,
AllowDynamic: &allowDynamic,
}
// Create the cAdvisor container manager
sysFs := sysfs.NewRealSysFs()
m, err := manager.New(memory.New(statsCacheDuration, nil), sysFs, housekeepingConfig, includedMetrics, http.DefaultClient, cgroupRoots, "")
if err != nil {
return nil, errors.Wrap(err, "new cadvisor manager")
}
if _, err := os.Stat(rootPath); err != nil {
if os.IsNotExist(err) {
if err := os.MkdirAll(path.Clean(rootPath), 0750); err != nil {
return nil, errors.Wrapf(err, "creating root direcotory %q", rootPath)
}
} else {
return nil, errors.Wrapf(err, "failed to stat %q", rootPath)
}
}
return &cadvisorClient{
imageFsInfoProvider: imageFsInfoProvider,
rootPath: rootPath,
Manager: m,
}, nil
}
func (cc *cadvisorClient) Start() error {
return cc.Manager.Start()
}
func (cc *cadvisorClient) ContainerInfo(name string, req *cadvisorapi.ContainerInfoRequest) (*cadvisorapi.ContainerInfo, error) {
return cc.GetContainerInfo(name, req)
}
func (cc *cadvisorClient) ContainerInfoV2(name string, options cadvisorapiv2.RequestOptions) (map[string]cadvisorapiv2.ContainerInfo, error) {
return cc.GetContainerInfoV2(name, options)
}
func (cc *cadvisorClient) VersionInfo() (*cadvisorapi.VersionInfo, error) {
return cc.GetVersionInfo()
}
func (cc *cadvisorClient) SubcontainerInfo(name string, req *cadvisorapi.ContainerInfoRequest) (map[string]*cadvisorapi.ContainerInfo, error) {
infos, err := cc.SubcontainersInfo(name, req)
if err != nil && len(infos) == 0 {
return nil, err
}
result := make(map[string]*cadvisorapi.ContainerInfo, len(infos))
for _, info := range infos {
result[info.Name] = info
}
return result, err
}
func (cc *cadvisorClient) MachineInfo() (*cadvisorapi.MachineInfo, error) {
return cc.GetMachineInfo()
}
func (cc *cadvisorClient) ImagesFsInfo() (cadvisorapiv2.FsInfo, error) {
label, err := cc.imageFsInfoProvider.ImageFsInfoLabel()
if err != nil {
return cadvisorapiv2.FsInfo{}, err
}
return cc.getFsInfo(label)
}
func (cc *cadvisorClient) RootFsInfo() (cadvisorapiv2.FsInfo, error) {
return cc.GetDirFsInfo(cc.rootPath)
}
func (cc *cadvisorClient) getFsInfo(label string) (cadvisorapiv2.FsInfo, error) {
res, err := cc.GetFsInfo(label)
if err != nil {
return cadvisorapiv2.FsInfo{}, err
}
if len(res) == 0 {
return cadvisorapiv2.FsInfo{}, fmt.Errorf("failed to find information for the filesystem labeled %q", label)
}
// TODO(vmarmol): Handle this better when a label has more than one image filesystem.
if len(res) > 1 {
log.Warningf("More than one filesystem labeled %q: %#v. Only using the first one", label, res)
}
return res[0], nil
}
func (cc *cadvisorClient) WatchEvents(request *events.Request) (*events.EventChannel, error) {
return cc.WatchForEvents(request)
}
+1
View File
@@ -0,0 +1 @@
package cadvisor // import "yunion.io/x/onecloud/pkg/util/pod/cadvisor"
+48
View File
@@ -0,0 +1,48 @@
package cadvisor
import (
cadvisorfs "github.com/google/cadvisor/fs"
"yunion.io/x/pkg/errors"
)
const (
DockerContainerRuntime = "docker"
RemoteContainerRuntime = "remote"
)
const (
// CrioSocket is the path to the CRI-O socket.
// Please keep this in sync with the one in:
// github.com/google/cadvisor/container/crio/client.go
CrioSocket = "/var/run/crio/crio.sock"
)
// imageFsInfoProvider knows how to translate the configured runtime
// to its file system label for images.
type imageFsInfoProvider struct {
runtime string
runtimeEndpoint string
}
func (i *imageFsInfoProvider) ImageFsInfoLabel() (string, error) {
switch i.runtime {
case DockerContainerRuntime:
return cadvisorfs.LabelDockerImages, nil
case RemoteContainerRuntime:
// This is a temporary workaround to get stats for cri-o from cadvisor
// and should be removed.
// Related to https://github.com/kubernetes/kubernetes/issues/51798
if i.runtimeEndpoint == CrioSocket || i.runtimeEndpoint == "unix://"+CrioSocket {
return cadvisorfs.LabelCrioImages, nil
}
}
return "", errors.Errorf("no imagefs label for configured runtime: %s", i.runtime)
}
func NewImageFsInfoProvider(runtime, endpoint string) ImageFsInfoProvider {
return &imageFsInfoProvider{
runtime: runtime,
runtimeEndpoint: endpoint,
}
}
+33
View File
@@ -0,0 +1,33 @@
package cadvisor
import (
"github.com/google/cadvisor/events"
cadvisorapi "github.com/google/cadvisor/info/v1"
cadvisorapiv2 "github.com/google/cadvisor/info/v2"
)
type Interface interface {
Start() error
ContainerInfo(name string, req *cadvisorapi.ContainerInfoRequest) (*cadvisorapi.ContainerInfo, error)
ContainerInfoV2(name string, options cadvisorapiv2.RequestOptions) (map[string]cadvisorapiv2.ContainerInfo, error)
MachineInfo() (*cadvisorapi.MachineInfo, error)
VersionInfo() (*cadvisorapi.VersionInfo, error)
// Returns usage information about the filesystem holding container images.
ImagesFsInfo() (cadvisorapiv2.FsInfo, error)
// Returns usage information about the root filesystem.
RootFsInfo() (cadvisorapiv2.FsInfo, error)
// Get events streamed through passedChannel that fit the request.
WatchEvents(request *events.Request) (*events.EventChannel, error)
// Get filesystem information for the filesystem that contains the given file.
GetDirFsInfo(path string) (cadvisorapiv2.FsInfo, error)
}
// ImageFsInfoProvider informs cAdvisor how to find imagefs for container images.
type ImageFsInfoProvider interface {
// ImageFsInfoLabel returns the label cAdvisor should use to find the filesystem holding container images.
ImageFsInfoLabel() (string, error)
}
@@ -0,0 +1,209 @@
package stats
import (
"fmt"
"path"
"sort"
"strings"
cadvisorapiv2 "github.com/google/cadvisor/info/v2"
"k8s.io/apimachinery/pkg/types"
"k8s.io/klog/v2"
"yunion.io/x/log"
"yunion.io/x/onecloud/pkg/util/pod/cadvisor"
)
// containerID is the identity of a container in a pod.
type containerID struct {
podRef PodReference
containerName string
}
// buildPodRef returns a PodReference that identifies the Pod managing cinfo
func buildPodRef(containerLabels map[string]string) PodReference {
podName := GetPodName(containerLabels)
podNamespace := GetPodNamespace(containerLabels)
podUID := GetPodUID(containerLabels)
return PodReference{Name: podName, Namespace: podNamespace, UID: podUID}
}
func getCadvisorContainerInfo(ca cadvisor.Interface) (map[string]cadvisorapiv2.ContainerInfo, error) {
infos, err := ca.ContainerInfoV2("/", cadvisorapiv2.RequestOptions{
IdType: cadvisorapiv2.TypeName,
Count: 2, // 2 samples are needed to compute "instantaneous" CPU
Recursive: true,
})
if err != nil {
if _, ok := infos["/"]; ok {
// If the failure is partial, log it and return a best-effort
// response.
log.Errorf("Partial failure issuing cadvisor.ContainerInfoV2: %v", err)
} else {
return nil, fmt.Errorf("failed to get root cgroup stats: %v", err)
}
}
return infos, nil
}
// libcontainerCgroupManagerType defines how to interface with libcontainer
type libcontainerCgroupManagerType string
const (
// libcontainerCgroupfs means use libcontainer with cgroupfs
libcontainerCgroupfs libcontainerCgroupManagerType = "cgroupfs"
// libcontainerSystemd means use libcontainer with systemd
libcontainerSystemd libcontainerCgroupManagerType = "systemd"
// systemdSuffix is the cgroup name suffix for systemd
systemdSuffix string = ".slice"
)
func IsSystemdStyleName(name string) bool {
return strings.HasSuffix(name, systemdSuffix)
}
// CgroupName is the abstract name of a cgroup prior to any driver specific conversion.
// It is specified as a list of strings from its individual components, such as:
// {"kubepods", "burstable", "pod1234-abcd-5678-efgh"}
type CgroupName []string
func ParseCgroupfsToCgroupName(name string) CgroupName {
components := strings.Split(strings.TrimPrefix(name, "/"), "/")
if len(components) == 1 && components[0] == "" {
components = []string{}
}
return CgroupName(components)
}
func unescapeSystemdCgroupName(part string) string {
return strings.Replace(part, "_", "-", -1)
}
func ParseSystemdToCgroupName(name string) CgroupName {
driverName := path.Base(name)
driverName = strings.TrimSuffix(driverName, systemdSuffix)
parts := strings.Split(driverName, "-")
result := []string{}
for _, part := range parts {
result = append(result, unescapeSystemdCgroupName(part))
}
return CgroupName(result)
}
const (
podCgroupNamePrefix = "pod"
)
// GetPodCgroupNameSuffix returns the last element of the pod CgroupName identifier
func GetPodCgroupNameSuffix(podUID types.UID) string {
return podCgroupNamePrefix + string(podUID)
}
// getCadvisorPodInfoFromPodUID returns a pod cgroup information by matching the podUID with its CgroupName identifier base name
func getCadvisorPodInfoFromPodUID(podUID types.UID, infos map[string]cadvisorapiv2.ContainerInfo) *cadvisorapiv2.ContainerInfo {
for key, info := range infos {
if IsSystemdStyleName(key) {
// Convert to internal cgroup name and take the last component only.
internalCgroupName := ParseSystemdToCgroupName(key)
key = internalCgroupName[len(internalCgroupName)-1]
} else {
// Take last component only.
key = path.Base(key)
}
if GetPodCgroupNameSuffix(podUID) == key {
return &info
}
}
return nil
}
// containerInfoWithCgroup contains the ContainerInfo and its cgroup name.
type containerInfoWithCgroup struct {
cinfo cadvisorapiv2.ContainerInfo
cgroup string
}
// removeTerminatedContainerInfo returns the specified containerInfo but with
// the stats of the terminated containers removed.
//
// A ContainerInfo is considered to be of a terminated container if it has an
// older CreationTime and zero CPU instantaneous and memory RSS usage.
func removeTerminatedContainerInfo(containerInfo map[string]cadvisorapiv2.ContainerInfo) map[string]cadvisorapiv2.ContainerInfo {
cinfoMap := make(map[containerID][]containerInfoWithCgroup)
for key, cinfo := range containerInfo {
if !isPodManagedContainer(&cinfo) {
continue
}
cinfoID := containerID{
podRef: buildPodRef(cinfo.Spec.Labels),
containerName: GetContainerName(cinfo.Spec.Labels),
}
cinfoMap[cinfoID] = append(cinfoMap[cinfoID], containerInfoWithCgroup{
cinfo: cinfo,
cgroup: key,
})
}
result := make(map[string]cadvisorapiv2.ContainerInfo)
for _, refs := range cinfoMap {
if len(refs) == 1 {
result[refs[0].cgroup] = refs[0].cinfo
continue
}
sort.Sort(ByCreationTime(refs))
for i := len(refs) - 1; i >= 0; i-- {
if hasMemoryAndCPUInstUsage(&refs[i].cinfo) {
result[refs[i].cgroup] = refs[i].cinfo
break
}
}
}
return result
}
// hasMemoryAndCPUInstUsage returns true if the specified container info has
// both non-zero CPU instantaneous usage and non-zero memory RSS usage, and
// false otherwise.
func hasMemoryAndCPUInstUsage(info *cadvisorapiv2.ContainerInfo) bool {
if !info.Spec.HasCpu || !info.Spec.HasMemory {
return false
}
cstat, found := latestContainerStats(info)
if !found {
return false
}
if cstat.CpuInst == nil {
return false
}
return cstat.CpuInst.Usage.Total != 0 && cstat.Memory.RSS != 0
}
// ByCreationTime implements sort.Interface for []containerInfoWithCgroup based
// on the cinfo.Spec.CreationTime field.
type ByCreationTime []containerInfoWithCgroup
func (a ByCreationTime) Len() int { return len(a) }
func (a ByCreationTime) Swap(i, j int) { a[i], a[j] = a[j], a[i] }
func (a ByCreationTime) Less(i, j int) bool {
if a[i].cinfo.Spec.CreationTime.Equal(a[j].cinfo.Spec.CreationTime) {
// There shouldn't be two containers with the same name and/or the same
// creation time. However, to make the logic here robust, we break the
// tie by moving the one without CPU instantaneous or memory RSS usage
// to the beginning.
return hasMemoryAndCPUInstUsage(&a[j].cinfo)
}
return a[i].cinfo.Spec.CreationTime.Before(a[j].cinfo.Spec.CreationTime)
}
// isPodManagedContainer returns true if the cinfo container is managed by a Pod
func isPodManagedContainer(cinfo *cadvisorapiv2.ContainerInfo) bool {
podName := GetPodName(cinfo.Spec.Labels)
podNamespace := GetPodNamespace(cinfo.Spec.Labels)
managed := podName != "" && podNamespace != ""
if !managed && podName != podNamespace {
klog.Warningf(
"Expect container to have either both podName (%s) and podNamespace (%s) labels, or neither.",
podName, podNamespace)
}
return managed
}
+805
View File
@@ -0,0 +1,805 @@
package stats
import (
"context"
"fmt"
"path"
"sort"
"strings"
"sync"
"time"
cadvisorfs "github.com/google/cadvisor/fs"
cadvisorapiv2 "github.com/google/cadvisor/info/v2"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/apimachinery/pkg/types"
runtimeapi "k8s.io/cri-api/pkg/apis/runtime/v1"
"k8s.io/klog/v2"
"yunion.io/x/pkg/errors"
"yunion.io/x/onecloud/pkg/util/pod/cadvisor"
)
var (
// defaultCachePeriod is the default cache period for each cpuUsage.
defaultCachePeriod = 10 * time.Minute
)
type cpuUsageRecord struct {
stats *runtimeapi.CpuUsage
usageNanoCores *uint64
}
// criStatsProvider implements the ContainerStatsProvider interface by getting
// the container stats from CRI.
type criStatsProvider struct {
// cadvisor is used to get the node root filesystem's stats (such as the
// capacity/available bytes/inodes) that will be populated in per container
// filesystem stats.
cadvisor cadvisor.Interface
// runtimeService is used to get the status and stats of the pods and its
// managed containers.
runtimeService runtimeapi.RuntimeServiceClient
// imageService is used to get the stats of the image filesystem.
imageService runtimeapi.ImageServiceClient
// cpuUsageCache caches the cpu usage for containers.
cpuUsageCache map[string]*cpuUsageRecord
mutex sync.RWMutex
}
func NewCRIContainerStatsProvider(
cadvisor cadvisor.Interface,
runtimeService runtimeapi.RuntimeServiceClient,
imageService runtimeapi.ImageServiceClient,
) ContainerStatsProvider {
return newCRIStatsProvider(cadvisor, runtimeService, imageService)
}
// newCRIStatsProvider returns a ContainerStatsProvider implementation that
// provides container stats using CRI.
func newCRIStatsProvider(
cadvisor cadvisor.Interface,
runtimeService runtimeapi.RuntimeServiceClient,
imageService runtimeapi.ImageServiceClient,
) ContainerStatsProvider {
return &criStatsProvider{
cadvisor: cadvisor,
runtimeService: runtimeService,
imageService: imageService,
cpuUsageCache: make(map[string]*cpuUsageRecord),
}
}
func (p *criStatsProvider) ListPodStats() ([]PodStats, error) {
// Don't update CPU nano core usage.
return p.listPodStats(false)
}
// ListPodStatsAndUpdateCPUNanoCoreUsage updates the cpu nano core usage for
// the containers and returns the stats for all the pod-managed containers.
// This is a workaround because CRI runtimes do not supply nano core usages,
// so this function calculate the difference between the current and the last
// (cached) cpu stats to calculate this metrics. The implementation assumes a
// single caller to periodically invoke this function to update the metrics. If
// there exist multiple callers, the period used to compute the cpu usage may
// vary and the usage could be incoherent (e.g., spiky). If no caller calls
// this function, the cpu usage will stay nil. Right now, eviction manager is
// the only caller, and it calls this function every 10s.
func (p *criStatsProvider) ListPodStatsAndUpdateCPUNanoCoreUsage() ([]PodStats, error) {
// Update CPU nano core usage.
return p.listPodStats(true)
}
func (p *criStatsProvider) listPodStats(updateCPUNanoCoreUsage bool) ([]PodStats, error) {
// Gets node root filesystem information, which will be used to populate
// the available and capacity bytes/inodes in container stats.
rootFsInfo, err := p.cadvisor.RootFsInfo()
if err != nil {
return nil, fmt.Errorf("failed to get rootFs info: %v", err)
}
csResp, err := p.runtimeService.ListContainers(context.Background(), &runtimeapi.ListContainersRequest{})
if err != nil {
return nil, errors.Wrap(err, "failed to list all containers")
}
containers := csResp.Containers
// Creates pod sandbox map.
podSandboxMap := make(map[string]*runtimeapi.PodSandbox)
resp, err := p.runtimeService.ListPodSandbox(context.Background(), &runtimeapi.ListPodSandboxRequest{})
if err != nil {
return nil, errors.Wrap(err, "failed to list all pod sandboxes")
}
podSandboxes := removeTerminatedPods(resp.Items)
for _, s := range podSandboxes {
podSandboxMap[s.Id] = s
}
// fsIDtoInfo is a map from filesystem id to its stats. This will be used
// as a cache to avoid querying cAdvisor for the filesystem stats with the
// same filesystem id many times.
fsIDtoInfo := make(map[runtimeapi.FilesystemIdentifier]*cadvisorapiv2.FsInfo)
// sandboxIDToPodStats is a temporary map from sandbox ID to its pod stats.
sandboxIDToPodStats := make(map[string]*PodStats)
cstsResp, err := p.runtimeService.ListContainerStats(context.Background(), &runtimeapi.ListContainerStatsRequest{})
if err != nil {
return nil, fmt.Errorf("failed to list all container stats: %v", err)
}
containers = removeTerminatedContainers(containers)
// Creates container map.
containerMap := make(map[string]*runtimeapi.Container)
for _, c := range containers {
containerMap[c.Id] = c
}
allInfos, err := getCadvisorContainerInfo(p.cadvisor)
if err != nil {
return nil, fmt.Errorf("failed to fetch cadvisor stats: %v", err)
}
caInfos := getCRICadvisorStats(allInfos)
// get network stats for containers.
// This is only used on Windows. For other platforms, (nil, nil) should be returned.
containerNetworkStats, err := p.listContainerNetworkStats()
if err != nil {
return nil, fmt.Errorf("failed to list container network stats: %v", err)
}
for _, stats := range cstsResp.Stats {
containerID := stats.Attributes.Id
container, found := containerMap[containerID]
if !found {
continue
}
podSandboxID := container.PodSandboxId
podSandbox, found := podSandboxMap[podSandboxID]
if !found {
continue
}
// Creates the stats of the pod (if not created yet) which the
// container belongs to.
ps, found := sandboxIDToPodStats[podSandboxID]
if !found {
ps = buildPodStats(podSandbox)
sandboxIDToPodStats[podSandboxID] = ps
}
// Fill available stats for full set of required pod stats
cs := p.makeContainerStats(stats, container, &rootFsInfo, fsIDtoInfo, podSandbox.GetMetadata(), updateCPUNanoCoreUsage)
p.addPodNetworkStats(ps, podSandboxID, caInfos, cs, containerNetworkStats[podSandboxID])
p.addPodCPUMemoryStats(ps, types.UID(podSandbox.Metadata.Uid), allInfos, cs)
p.addProcessStats(ps, types.UID(podSandbox.Metadata.Uid), allInfos, cs)
// If cadvisor stats is available for the container, use it to populate
// container stats
caStats, caFound := caInfos[containerID]
if !caFound {
klog.V(5).Infof("Unable to find cadvisor stats for %q", containerID)
} else {
p.addCadvisorContainerStats(cs, &caStats)
}
ps.Containers = append(ps.Containers, *cs)
}
// cleanup outdated caches.
p.cleanupOutdatedCaches()
result := make([]PodStats, 0, len(sandboxIDToPodStats))
for _, s := range sandboxIDToPodStats {
//p.makePodStorageStats(s, &rootFsInfo)
result = append(result, *s)
}
return result, nil
}
func (p *criStatsProvider) ListPodCPUAndMemoryStats() ([]PodStats, error) {
ctx := context.Background()
containersResp, err := p.runtimeService.ListContainers(ctx, &runtimeapi.ListContainersRequest{})
if err != nil {
return nil, fmt.Errorf("failed to list all containers: %v", err)
}
containers := containersResp.Containers
// Creates pod sandbox map.
podSandboxMap := make(map[string]*runtimeapi.PodSandbox)
resp, err := p.runtimeService.ListPodSandbox(ctx, &runtimeapi.ListPodSandboxRequest{})
if err != nil {
return nil, fmt.Errorf("failed to list all pod sandboxes: %v", err)
}
podSandboxes := resp.Items
podSandboxes = removeTerminatedPods(podSandboxes)
for _, s := range podSandboxes {
podSandboxMap[s.Id] = s
}
// sandboxIDToPodStats is a temporary map from sandbox ID to its pod stats.
sandboxIDToPodStats := make(map[string]*PodStats)
containerStatResp, err := p.runtimeService.ListContainerStats(ctx, &runtimeapi.ListContainerStatsRequest{})
if err != nil {
return nil, fmt.Errorf("failed to list all container stats: %v", err)
}
containers = removeTerminatedContainers(containers)
// Creates container map.
containerMap := make(map[string]*runtimeapi.Container)
for _, c := range containers {
containerMap[c.Id] = c
}
allInfos, err := getCadvisorContainerInfo(p.cadvisor)
if err != nil {
return nil, fmt.Errorf("failed to fetch cadvisor stats: %v", err)
}
caInfos := getCRICadvisorStats(allInfos)
for _, stats := range containerStatResp.Stats {
containerID := stats.Attributes.Id
container, found := containerMap[containerID]
if !found {
continue
}
podSandboxID := container.PodSandboxId
podSandbox, found := podSandboxMap[podSandboxID]
if !found {
continue
}
// Creates the stats of the pod (if not created yet) which the
// container belongs to.
ps, found := sandboxIDToPodStats[podSandboxID]
if !found {
ps = buildPodStats(podSandbox)
sandboxIDToPodStats[podSandboxID] = ps
}
// Fill available CPU and memory stats for full set of required pod stats
cs := p.makeContainerCPUAndMemoryStats(stats, container)
p.addPodCPUMemoryStats(ps, types.UID(podSandbox.Metadata.Uid), allInfos, cs)
// If cadvisor stats is available for the container, use it to populate
// container stats
caStats, caFound := caInfos[containerID]
if !caFound {
klog.V(4).Infof("Unable to find cadvisor stats for %q", containerID)
} else {
p.addCadvisorContainerStats(cs, &caStats)
}
ps.Containers = append(ps.Containers, *cs)
}
// cleanup outdated caches.
p.cleanupOutdatedCaches()
result := make([]PodStats, 0, len(sandboxIDToPodStats))
for _, s := range sandboxIDToPodStats {
result = append(result, *s)
}
return result, nil
}
func (p *criStatsProvider) ImageFsStats() (FsStats, error) {
//TODO implement me
panic("implement me")
}
func (p *criStatsProvider) ImageFsDevice() (string, error) {
//TODO implement me
panic("implement me")
}
// buildPodStats returns a PodStats that identifies the Pod managing cinfo
func buildPodStats(podSandbox *runtimeapi.PodSandbox) *PodStats {
return &PodStats{
PodRef: PodReference{
Name: podSandbox.Metadata.Name,
UID: podSandbox.Metadata.Uid,
Namespace: podSandbox.Metadata.Namespace,
},
// The StartTime in the summary API is the pod creation time.
StartTime: metav1.NewTime(time.Unix(0, podSandbox.CreatedAt)),
}
}
/*func (p *criStatsProvider) makePodStorageStats(s *PodStats, rootFsInfo *cadvisorapiv2.FsInfo) {
podNs := s.PodRef.Namespace
podName := s.PodRef.Name
podUID := types.UID(s.PodRef.UID)
vstats, found := p.resourceAnalyzer.GetPodVolumeStats(podUID)
if !found {
return
}
podLogDir := kuberuntime.BuildPodLogsDirectory(podNs, podName, podUID)
logStats, err := p.getPodLogStats(podLogDir, rootFsInfo)
if err != nil {
klog.Errorf("Unable to fetch pod log stats for path %s: %v ", podLogDir, err)
// If people do in-place upgrade, there might be pods still using
// the old log path. For those pods, no pod log stats is returned.
// We should continue generating other stats in that case.
// calcEphemeralStorage tolerants logStats == nil.
}
ephemeralStats := make([]statsapi.VolumeStats, len(vstats.EphemeralVolumes))
copy(ephemeralStats, vstats.EphemeralVolumes)
s.VolumeStats = append(append([]statsapi.VolumeStats{}, vstats.EphemeralVolumes...), vstats.PersistentVolumes...)
s.EphemeralStorage = calcEphemeralStorage(s.Containers, ephemeralStats, rootFsInfo, logStats, true)
}*/
func (p *criStatsProvider) addPodNetworkStats(
ps *PodStats,
podSandboxID string,
caInfos map[string]cadvisorapiv2.ContainerInfo,
cs *ContainerStats,
netStats *NetworkStats,
) {
caPodSandbox, found := caInfos[podSandboxID]
// try get network stats from cadvisor first.
if found {
networkStats := cadvisorInfoToNetworkStats(&caPodSandbox)
if networkStats != nil {
ps.Network = networkStats
return
}
}
// Not found from cadvisor, get from netStats.
if netStats != nil {
ps.Network = netStats
return
}
// TODO: sum Pod network stats from container stats.
klog.V(4).Infof("Unable to find network stats for sandbox %q", podSandboxID)
}
func (p *criStatsProvider) addPodCPUMemoryStats(
ps *PodStats,
podUID types.UID,
allInfos map[string]cadvisorapiv2.ContainerInfo,
cs *ContainerStats,
) {
// try get cpu and memory stats from cadvisor first.
podCgroupInfo := getCadvisorPodInfoFromPodUID(podUID, allInfos)
if podCgroupInfo != nil {
cpu, memory := cadvisorInfoToCPUandMemoryStats(podCgroupInfo)
ps.CPU = cpu
ps.Memory = memory
return
}
// Sum Pod cpu and memory stats from containers stats.
if cs.CPU != nil {
if ps.CPU == nil {
ps.CPU = &CPUStats{}
}
ps.CPU.Time = cs.CPU.Time
usageCoreNanoSeconds := getUint64Value(cs.CPU.UsageCoreNanoSeconds) + getUint64Value(ps.CPU.UsageCoreNanoSeconds)
usageNanoCores := getUint64Value(cs.CPU.UsageNanoCores) + getUint64Value(ps.CPU.UsageNanoCores)
ps.CPU.UsageCoreNanoSeconds = &usageCoreNanoSeconds
ps.CPU.UsageNanoCores = &usageNanoCores
}
if cs.Memory != nil {
if ps.Memory == nil {
ps.Memory = &MemoryStats{}
}
ps.Memory.Time = cs.Memory.Time
availableBytes := getUint64Value(cs.Memory.AvailableBytes) + getUint64Value(ps.Memory.AvailableBytes)
usageBytes := getUint64Value(cs.Memory.UsageBytes) + getUint64Value(ps.Memory.UsageBytes)
workingSetBytes := getUint64Value(cs.Memory.WorkingSetBytes) + getUint64Value(ps.Memory.WorkingSetBytes)
rSSBytes := getUint64Value(cs.Memory.RSSBytes) + getUint64Value(ps.Memory.RSSBytes)
pageFaults := getUint64Value(cs.Memory.PageFaults) + getUint64Value(ps.Memory.PageFaults)
majorPageFaults := getUint64Value(cs.Memory.MajorPageFaults) + getUint64Value(ps.Memory.MajorPageFaults)
ps.Memory.AvailableBytes = &availableBytes
ps.Memory.UsageBytes = &usageBytes
ps.Memory.WorkingSetBytes = &workingSetBytes
ps.Memory.RSSBytes = &rSSBytes
ps.Memory.PageFaults = &pageFaults
ps.Memory.MajorPageFaults = &majorPageFaults
}
}
func (p *criStatsProvider) addProcessStats(
ps *PodStats,
podUID types.UID,
allInfos map[string]cadvisorapiv2.ContainerInfo,
cs *ContainerStats,
) {
// try get process stats from cadvisor only.
info := getCadvisorPodInfoFromPodUID(podUID, allInfos)
if info != nil {
ps.ProcessStats = cadvisorInfoToProcessStats(info)
return
}
}
// getFsInfo returns the information of the filesystem with the specified
// fsID. If any error occurs, this function logs the error and returns
// nil.
func (p *criStatsProvider) getFsInfo(fsID *runtimeapi.FilesystemIdentifier) *cadvisorapiv2.FsInfo {
if fsID == nil {
klog.V(2).Infof("Failed to get filesystem info: fsID is nil.")
return nil
}
mountpoint := fsID.GetMountpoint()
fsInfo, err := p.cadvisor.GetDirFsInfo(mountpoint)
if err != nil {
msg := fmt.Sprintf("Failed to get the info of the filesystem with mountpoint %q: %v.", mountpoint, err)
if err == cadvisorfs.ErrNoSuchDevice {
klog.V(2).Info(msg)
} else {
klog.Error(msg)
}
return nil
}
return &fsInfo
}
func (p *criStatsProvider) makeContainerStats(
stats *runtimeapi.ContainerStats,
container *runtimeapi.Container,
rootFsInfo *cadvisorapiv2.FsInfo,
fsIDtoInfo map[runtimeapi.FilesystemIdentifier]*cadvisorapiv2.FsInfo,
meta *runtimeapi.PodSandboxMetadata,
updateCPUNanoCoreUsage bool,
) *ContainerStats {
result := &ContainerStats{
Name: stats.Attributes.Metadata.Name,
// The StartTime in the summary API is the container creation time.
StartTime: metav1.NewTime(time.Unix(0, container.CreatedAt)),
CPU: &CPUStats{},
Memory: &MemoryStats{},
Rootfs: &FsStats{},
// UserDefinedMetrics is not supported by CRI.
}
if stats.Cpu != nil {
result.CPU.Time = metav1.NewTime(time.Unix(0, stats.Cpu.Timestamp))
if stats.Cpu.UsageCoreNanoSeconds != nil {
result.CPU.UsageCoreNanoSeconds = &stats.Cpu.UsageCoreNanoSeconds.Value
}
var usageNanoCores *uint64
if updateCPUNanoCoreUsage {
usageNanoCores = p.getAndUpdateContainerUsageNanoCores(stats)
} else {
usageNanoCores = p.getContainerUsageNanoCores(stats)
}
if usageNanoCores != nil {
result.CPU.UsageNanoCores = usageNanoCores
}
} else {
result.CPU.Time = metav1.NewTime(time.Unix(0, time.Now().UnixNano()))
result.CPU.UsageCoreNanoSeconds = uint64Ptr(0)
result.CPU.UsageNanoCores = uint64Ptr(0)
}
if stats.Memory != nil {
result.Memory.Time = metav1.NewTime(time.Unix(0, stats.Memory.Timestamp))
if stats.Memory.WorkingSetBytes != nil {
result.Memory.WorkingSetBytes = &stats.Memory.WorkingSetBytes.Value
}
} else {
result.Memory.Time = metav1.NewTime(time.Unix(0, time.Now().UnixNano()))
result.Memory.WorkingSetBytes = uint64Ptr(0)
}
if stats.WritableLayer != nil {
result.Rootfs.Time = metav1.NewTime(time.Unix(0, stats.WritableLayer.Timestamp))
if stats.WritableLayer.UsedBytes != nil {
result.Rootfs.UsedBytes = &stats.WritableLayer.UsedBytes.Value
}
if stats.WritableLayer.InodesUsed != nil {
result.Rootfs.InodesUsed = &stats.WritableLayer.InodesUsed.Value
}
}
fsID := stats.GetWritableLayer().GetFsId()
if fsID != nil {
imageFsInfo, found := fsIDtoInfo[*fsID]
if !found {
imageFsInfo = p.getFsInfo(fsID)
fsIDtoInfo[*fsID] = imageFsInfo
}
if imageFsInfo != nil {
// The image filesystem id is unknown to the local node or there's
// an error on retrieving the stats. In these cases, we omit those stats
// and return the best-effort partial result. See
// https://github.com/kubernetes/heapster/issues/1793.
result.Rootfs.AvailableBytes = &imageFsInfo.Available
result.Rootfs.CapacityBytes = &imageFsInfo.Capacity
result.Rootfs.InodesFree = imageFsInfo.InodesFree
result.Rootfs.Inodes = imageFsInfo.Inodes
}
}
// NOTE: This doesn't support the old pod log path, `/var/log/pods/UID`. For containers
// using old log path, empty log stats are returned. This is fine, because we don't
// officially support in-place upgrade anyway.
/*var (
containerLogPath = kuberuntime.BuildContainerLogsDirectory(meta.GetNamespace(),
meta.GetName(), types.UID(meta.GetUid()), container.GetMetadata().GetName())
err error
)
result.Logs, err = p.getPathFsStats(containerLogPath, rootFsInfo)
if err != nil {
klog.Errorf("Unable to fetch container log stats for path %s: %v ", containerLogPath, err)
}*/
return result
}
func (p *criStatsProvider) makeContainerCPUAndMemoryStats(
stats *runtimeapi.ContainerStats,
container *runtimeapi.Container,
) *ContainerStats {
result := &ContainerStats{
Name: stats.Attributes.Metadata.Name,
// The StartTime in the summary API is the container creation time.
StartTime: metav1.NewTime(time.Unix(0, container.CreatedAt)),
CPU: &CPUStats{},
Memory: &MemoryStats{},
// UserDefinedMetrics is not supported by CRI.
}
if stats.Cpu != nil {
result.CPU.Time = metav1.NewTime(time.Unix(0, stats.Cpu.Timestamp))
if stats.Cpu.UsageCoreNanoSeconds != nil {
result.CPU.UsageCoreNanoSeconds = &stats.Cpu.UsageCoreNanoSeconds.Value
}
usageNanoCores := p.getContainerUsageNanoCores(stats)
if usageNanoCores != nil {
result.CPU.UsageNanoCores = usageNanoCores
}
} else {
result.CPU.Time = metav1.NewTime(time.Unix(0, time.Now().UnixNano()))
result.CPU.UsageCoreNanoSeconds = uint64Ptr(0)
result.CPU.UsageNanoCores = uint64Ptr(0)
}
if stats.Memory != nil {
result.Memory.Time = metav1.NewTime(time.Unix(0, stats.Memory.Timestamp))
if stats.Memory.WorkingSetBytes != nil {
result.Memory.WorkingSetBytes = &stats.Memory.WorkingSetBytes.Value
}
} else {
result.Memory.Time = metav1.NewTime(time.Unix(0, time.Now().UnixNano()))
result.Memory.WorkingSetBytes = uint64Ptr(0)
}
return result
}
// getContainerUsageNanoCores gets the cached usageNanoCores.
func (p *criStatsProvider) getContainerUsageNanoCores(stats *runtimeapi.ContainerStats) *uint64 {
if stats == nil || stats.Attributes == nil {
return nil
}
p.mutex.RLock()
defer p.mutex.RUnlock()
cached, ok := p.cpuUsageCache[stats.Attributes.Id]
if !ok || cached.usageNanoCores == nil {
return nil
}
// return a copy of the usage
latestUsage := *cached.usageNanoCores
return &latestUsage
}
// getContainerUsageNanoCores computes usageNanoCores based on the given and
// the cached usageCoreNanoSeconds, updates the cache with the computed
// usageNanoCores, and returns the usageNanoCores.
func (p *criStatsProvider) getAndUpdateContainerUsageNanoCores(stats *runtimeapi.ContainerStats) *uint64 {
if stats == nil || stats.Attributes == nil || stats.Cpu == nil || stats.Cpu.UsageCoreNanoSeconds == nil {
return nil
}
id := stats.Attributes.Id
usage, err := func() (*uint64, error) {
p.mutex.Lock()
defer p.mutex.Unlock()
cached, ok := p.cpuUsageCache[id]
if !ok || cached.stats.UsageCoreNanoSeconds == nil || stats.Cpu.UsageCoreNanoSeconds.Value < cached.stats.UsageCoreNanoSeconds.Value {
// Cannot compute the usage now, but update the cached stats anyway
p.cpuUsageCache[id] = &cpuUsageRecord{stats: stats.Cpu, usageNanoCores: nil}
return nil, nil
}
newStats := stats.Cpu
cachedStats := cached.stats
nanoSeconds := newStats.Timestamp - cachedStats.Timestamp
if nanoSeconds <= 0 {
return nil, fmt.Errorf("zero or negative interval (%v - %v)", newStats.Timestamp, cachedStats.Timestamp)
}
usageNanoCores := uint64(float64(newStats.UsageCoreNanoSeconds.Value-cachedStats.UsageCoreNanoSeconds.Value) /
float64(nanoSeconds) * float64(time.Second/time.Nanosecond))
// Update cache with new value.
usageToUpdate := usageNanoCores
p.cpuUsageCache[id] = &cpuUsageRecord{stats: newStats, usageNanoCores: &usageToUpdate}
return &usageNanoCores, nil
}()
if err != nil {
// This should not happen. Log now to raise visibility
klog.Errorf("failed updating cpu usage nano core: %v", err)
}
return usage
}
func (p *criStatsProvider) cleanupOutdatedCaches() {
p.mutex.Lock()
defer p.mutex.Unlock()
for k, v := range p.cpuUsageCache {
if v == nil {
delete(p.cpuUsageCache, k)
continue
}
if time.Since(time.Unix(0, v.stats.Timestamp)) > defaultCachePeriod {
delete(p.cpuUsageCache, k)
}
}
}
// removeTerminatedPods returns pods with terminated ones removed.
// It only removes a terminated pod when there is a running instance
// of the pod with the same name and namespace.
// This is needed because:
// 1) PodSandbox may be recreated;
// 2) Pod may be recreated with the same name and namespace.
func removeTerminatedPods(pods []*runtimeapi.PodSandbox) []*runtimeapi.PodSandbox {
podMap := make(map[PodReference][]*runtimeapi.PodSandbox)
// Sort order by create time
sort.Slice(pods, func(i, j int) bool {
return pods[i].CreatedAt < pods[j].CreatedAt
})
for _, pod := range pods {
refID := PodReference{
Name: pod.GetMetadata().GetName(),
Namespace: pod.GetMetadata().GetNamespace(),
// UID is intentionally left empty.
}
podMap[refID] = append(podMap[refID], pod)
}
result := make([]*runtimeapi.PodSandbox, 0)
for _, refs := range podMap {
if len(refs) == 1 {
result = append(result, refs[0])
continue
}
found := false
for i := 0; i < len(refs); i++ {
if refs[i].State == runtimeapi.PodSandboxState_SANDBOX_READY {
found = true
result = append(result, refs[i])
}
}
if !found {
result = append(result, refs[len(refs)-1])
}
}
return result
}
// removeTerminatedContainers removes all terminated containers since they should
// not be used for usage calculations.
func removeTerminatedContainers(containers []*runtimeapi.Container) []*runtimeapi.Container {
containerMap := make(map[containerID][]*runtimeapi.Container)
// Sort order by create time
sort.Slice(containers, func(i, j int) bool {
return containers[i].CreatedAt < containers[j].CreatedAt
})
for _, container := range containers {
refID := containerID{
podRef: buildPodRef(container.Labels),
containerName: GetContainerName(container.Labels),
}
containerMap[refID] = append(containerMap[refID], container)
}
result := make([]*runtimeapi.Container, 0)
for _, refs := range containerMap {
for i := 0; i < len(refs); i++ {
if refs[i].State == runtimeapi.ContainerState_CONTAINER_RUNNING {
result = append(result, refs[i])
}
}
}
return result
}
func (p *criStatsProvider) addCadvisorContainerStats(
cs *ContainerStats,
caPodStats *cadvisorapiv2.ContainerInfo,
) {
if caPodStats.Spec.HasCustomMetrics {
cs.UserDefinedMetrics = cadvisorInfoToUserDefinedMetrics(caPodStats)
}
cpu, memory := cadvisorInfoToCPUandMemoryStats(caPodStats)
if cpu != nil {
cs.CPU = cpu
}
if memory != nil {
cs.Memory = memory
}
}
func getCRICadvisorStats(infos map[string]cadvisorapiv2.ContainerInfo) map[string]cadvisorapiv2.ContainerInfo {
stats := make(map[string]cadvisorapiv2.ContainerInfo)
infos = removeTerminatedContainerInfo(infos)
for key, info := range infos {
// On systemd using devicemapper each mount into the container has an
// associated cgroup. We ignore them to ensure we do not get duplicate
// entries in our summary. For details on .mount units:
// http://man7.org/linux/man-pages/man5/systemd.mount.5.html
if strings.HasSuffix(key, ".mount") {
continue
}
// Build the Pod key if this container is managed by a Pod
if !isPodManagedContainer(&info) {
continue
}
stats[path.Base(key)] = info
}
return stats
}
/*func (p *criStatsProvider) getPathFsStats(path string, rootFsInfo *cadvisorapiv2.FsInfo) (*statsapi.FsStats, error) {
m := p.logMetricsService.createLogMetricsProvider(path)
logMetrics, err := m.GetMetrics()
if err != nil {
return nil, err
}
result := &statsapi.FsStats{
Time: metav1.NewTime(rootFsInfo.Timestamp),
AvailableBytes: &rootFsInfo.Available,
CapacityBytes: &rootFsInfo.Capacity,
InodesFree: rootFsInfo.InodesFree,
Inodes: rootFsInfo.Inodes,
}
usedbytes := uint64(logMetrics.Used.Value())
result.UsedBytes = &usedbytes
inodesUsed := uint64(logMetrics.InodesUsed.Value())
result.InodesUsed = &inodesUsed
result.Time = maxUpdateTime(&result.Time, &logMetrics.Time)
return result, nil
}*/
// getPodLogStats gets stats for logs under the pod log directory. Container logs usually exist
// under the container log directory. However, for some container runtimes, e.g. kata, gvisor,
// they may want to keep some pod level logs, in that case they can put those logs directly under
// the pod log directory. And kubelet will take those logs into account as part of pod ephemeral
// storage.
/*func (p *criStatsProvider) getPodLogStats(path string, rootFsInfo *cadvisorapiv2.FsInfo) (*statsapi.FsStats, error) {
files, err := p.osInterface.ReadDir(path)
if err != nil {
return nil, err
}
result := &statsapi.FsStats{
Time: metav1.NewTime(rootFsInfo.Timestamp),
AvailableBytes: &rootFsInfo.Available,
CapacityBytes: &rootFsInfo.Capacity,
InodesFree: rootFsInfo.InodesFree,
Inodes: rootFsInfo.Inodes,
}
for _, f := range files {
if f.IsDir() {
continue
}
// Only include *files* under pod log directory.
fpath := filepath.Join(path, f.Name())
fstats, err := p.getPathFsStats(fpath, rootFsInfo)
if err != nil {
return nil, fmt.Errorf("failed to get fsstats for %q: %v", fpath, err)
}
result.UsedBytes = addUsage(result.UsedBytes, fstats.UsedBytes)
result.InodesUsed = addUsage(result.InodesUsed, fstats.InodesUsed)
result.Time = maxUpdateTime(&result.Time, &fstats.Time)
}
return result, nil
}*/
@@ -0,0 +1,7 @@
package stats
// listContainerNetworkStats returns the network stats of all the running containers.
// It should return (nil, nil) for platforms other than Windows.
func (p *criStatsProvider) listContainerNetworkStats() (map[string]*NetworkStats, error) {
return nil, nil
}
+1
View File
@@ -0,0 +1 @@
package stats // import "yunion.io/x/onecloud/pkg/util/pod/stats"
+193
View File
@@ -0,0 +1,193 @@
package stats
import (
"time"
cadvisorapiv1 "github.com/google/cadvisor/info/v1"
cadvisorapiv2 "github.com/google/cadvisor/info/v2"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/klog/v2"
)
// defaultNetworkInterfaceName is used for collectng network stats.
// This logic relies on knowledge of the container runtime implementation and
// is not reliable.
const defaultNetworkInterfaceName = "eth0"
func getUint64Value(value *uint64) uint64 {
if value == nil {
return 0
}
return *value
}
func uint64Ptr(i uint64) *uint64 {
return &i
}
func cadvisorInfoToCPUandMemoryStats(info *cadvisorapiv2.ContainerInfo) (*CPUStats, *MemoryStats) {
cstat, found := latestContainerStats(info)
if !found {
return nil, nil
}
var cpuStats *CPUStats
var memoryStats *MemoryStats
cpuStats = &CPUStats{
Time: metav1.NewTime(cstat.Timestamp),
UsageNanoCores: uint64Ptr(0),
UsageCoreNanoSeconds: uint64Ptr(0),
}
if info.Spec.HasCpu {
if cstat.CpuInst != nil {
cpuStats.UsageNanoCores = &cstat.CpuInst.Usage.Total
}
if cstat.Cpu != nil {
cpuStats.UsageCoreNanoSeconds = &cstat.Cpu.Usage.Total
}
}
if info.Spec.HasMemory && cstat.Memory != nil {
pageFaults := cstat.Memory.ContainerData.Pgfault
majorPageFaults := cstat.Memory.ContainerData.Pgmajfault
memoryStats = &MemoryStats{
Time: metav1.NewTime(cstat.Timestamp),
UsageBytes: &cstat.Memory.Usage,
WorkingSetBytes: &cstat.Memory.WorkingSet,
RSSBytes: &cstat.Memory.RSS,
PageFaults: &pageFaults,
MajorPageFaults: &majorPageFaults,
}
// availableBytes = memory limit (if known) - workingset
if !isMemoryUnlimited(info.Spec.Memory.Limit) {
availableBytes := info.Spec.Memory.Limit - cstat.Memory.WorkingSet
memoryStats.AvailableBytes = &availableBytes
}
} else {
memoryStats = &MemoryStats{
Time: metav1.NewTime(cstat.Timestamp),
WorkingSetBytes: uint64Ptr(0),
}
}
return cpuStats, memoryStats
}
// latestContainerStats returns the latest container stats from cadvisor, or nil if none exist
func latestContainerStats(info *cadvisorapiv2.ContainerInfo) (*cadvisorapiv2.ContainerStats, bool) {
stats := info.Stats
if len(stats) < 1 {
return nil, false
}
latest := stats[len(stats)-1]
if latest == nil {
return nil, false
}
return latest, true
}
func isMemoryUnlimited(v uint64) bool {
// Size after which we consider memory to be "unlimited". This is not
// MaxInt64 due to rounding by the kernel.
// TODO: cadvisor should export this https://github.com/google/cadvisor/blob/master/metrics/prometheus.go#L596
const maxMemorySize = uint64(1 << 62)
return v > maxMemorySize
}
// cadvisorInfoToNetworkStats returns the statsapi.NetworkStats converted from
// the container info from cadvisor.
func cadvisorInfoToNetworkStats(info *cadvisorapiv2.ContainerInfo) *NetworkStats {
if !info.Spec.HasNetwork {
return nil
}
cstat, found := latestContainerStats(info)
if !found {
return nil
}
if cstat.Network == nil {
return nil
}
iStats := NetworkStats{
Time: metav1.NewTime(cstat.Timestamp),
}
for i := range cstat.Network.Interfaces {
inter := cstat.Network.Interfaces[i]
iStat := InterfaceStats{
Name: inter.Name,
RxBytes: &inter.RxBytes,
RxErrors: &inter.RxErrors,
TxBytes: &inter.TxBytes,
TxErrors: &inter.TxErrors,
}
if inter.Name == defaultNetworkInterfaceName {
iStats.InterfaceStats = iStat
}
iStats.Interfaces = append(iStats.Interfaces, iStat)
}
return &iStats
}
// cadvisorInfoToUserDefinedMetrics returns the statsapi.UserDefinedMetric
// converted from the container info from cadvisor.
func cadvisorInfoToUserDefinedMetrics(info *cadvisorapiv2.ContainerInfo) []UserDefinedMetric {
type specVal struct {
ref UserDefinedMetricDescriptor
valType cadvisorapiv1.DataType
time time.Time
value float64
}
udmMap := map[string]*specVal{}
for _, spec := range info.Spec.CustomMetrics {
udmMap[spec.Name] = &specVal{
ref: UserDefinedMetricDescriptor{
Name: spec.Name,
Type: UserDefinedMetricType(spec.Type),
Units: spec.Units,
},
valType: spec.Format,
}
}
for _, stat := range info.Stats {
for name, values := range stat.CustomMetrics {
specVal, ok := udmMap[name]
if !ok {
klog.Warningf("spec for custom metric %q is missing from cAdvisor output. Spec: %+v, Metrics: %+v", name, info.Spec, stat.CustomMetrics)
continue
}
for _, value := range values {
// Pick the most recent value
if value.Timestamp.Before(specVal.time) {
continue
}
specVal.time = value.Timestamp
specVal.value = value.FloatValue
if specVal.valType == cadvisorapiv1.IntType {
specVal.value = float64(value.IntValue)
}
}
}
}
var udm []UserDefinedMetric
for _, specVal := range udmMap {
udm = append(udm, UserDefinedMetric{
UserDefinedMetricDescriptor: specVal.ref,
Time: metav1.NewTime(specVal.time),
Value: specVal.value,
})
}
return udm
}
func cadvisorInfoToProcessStats(info *cadvisorapiv2.ContainerInfo) *ProcessStats {
cstat, found := latestContainerStats(info)
if !found || cstat.Processes == nil {
return nil
}
num := cstat.Processes.ProcessCount
return &ProcessStats{ProcessCount: uint64Ptr(num)}
}
+30
View File
@@ -0,0 +1,30 @@
package stats
import "yunion.io/x/onecloud/pkg/util/pod/cadvisor"
type ContainerStatsProvider interface {
ListPodStats() ([]PodStats, error)
ListPodStatsAndUpdateCPUNanoCoreUsage() ([]PodStats, error)
ListPodCPUAndMemoryStats() ([]PodStats, error)
ImageFsStats() (FsStats, error)
ImageFsDevice() (string, error)
}
type StatsProvider struct {
cadvisor cadvisor.Interface
ContainerStatsProvider
}
func NewCRIStatsProvider(
cadvisor cadvisor.Interface,
) *StatsProvider {
return nil
}
func newStatsProvider(
cadvisor cadvisor.Interface,
) *StatsProvider {
return &StatsProvider{
cadvisor: cadvisor,
}
}
+350
View File
@@ -0,0 +1,350 @@
package stats
import metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
// Summary is a top-level container for holding NodeStats and PodStats.
type Summary struct {
// Overall node stats.
Node NodeStats `json:"node"`
// Per-pod stats.
Pods []PodStats `json:"pods"`
}
// NodeStats holds node-level unprocessed sample stats.
type NodeStats struct {
// Reference to the measured Node.
NodeName string `json:"nodeName"`
// Stats of system daemons tracked as raw containers.
// The system containers are named according to the SystemContainer* constants.
// +optional
// +patchMergeKey=name
// +patchStrategy=merge
SystemContainers []ContainerStats `json:"systemContainers,omitempty" patchStrategy:"merge" patchMergeKey:"name"`
// The time at which data collection for the node-scoped (i.e. aggregate) stats was (re)started.
StartTime metav1.Time `json:"startTime"`
// Stats pertaining to CPU resources.
// +optional
CPU *CPUStats `json:"cpu,omitempty"`
// Stats pertaining to memory (RAM) resources.
// +optional
Memory *MemoryStats `json:"memory,omitempty"`
// Stats pertaining to network resources.
// +optional
Network *NetworkStats `json:"network,omitempty"`
// Stats pertaining to total usage of filesystem resources on the rootfs used by node k8s components.
// NodeFs.Used is the total bytes used on the filesystem.
// +optional
Fs *FsStats `json:"fs,omitempty"`
// Stats about the underlying container runtime.
// +optional
Runtime *RuntimeStats `json:"runtime,omitempty"`
// Stats about the rlimit of system.
// +optional
Rlimit *RlimitStats `json:"rlimit,omitempty"`
}
// RlimitStats are stats rlimit of OS.
type RlimitStats struct {
Time metav1.Time `json:"time"`
// The max PID of OS.
MaxPID *int64 `json:"maxpid,omitempty"`
// The number of running process in the OS.
NumOfRunningProcesses *int64 `json:"curproc,omitempty"`
}
// RuntimeStats are stats pertaining to the underlying container runtime.
type RuntimeStats struct {
// Stats about the underlying filesystem where container images are stored.
// This filesystem could be the same as the primary (root) filesystem.
// Usage here refers to the total number of bytes occupied by images on the filesystem.
// +optional
ImageFs *FsStats `json:"imageFs,omitempty"`
}
const (
// SystemContainerKubelet is the container name for the system container tracking Kubelet usage.
SystemContainerKubelet = "kubelet"
// SystemContainerRuntime is the container name for the system container tracking the runtime (e.g. docker) usage.
SystemContainerRuntime = "runtime"
// SystemContainerMisc is the container name for the system container tracking non-kubernetes processes.
SystemContainerMisc = "misc"
// SystemContainerPods is the container name for the system container tracking user pods.
SystemContainerPods = "pods"
)
// PodStats holds pod-level unprocessed sample stats.
type PodStats struct {
// Reference to the measured Pod.
PodRef PodReference `json:"podRef"`
// The time at which data collection for the pod-scoped (e.g. network) stats was (re)started.
StartTime metav1.Time `json:"startTime"`
// Stats of containers in the measured pod.
// +patchMergeKey=name
// +patchStrategy=merge
Containers []ContainerStats `json:"containers" patchStrategy:"merge" patchMergeKey:"name"`
// Stats pertaining to CPU resources consumed by pod cgroup (which includes all containers' resource usage and pod overhead).
// +optional
CPU *CPUStats `json:"cpu,omitempty"`
// Stats pertaining to memory (RAM) resources consumed by pod cgroup (which includes all containers' resource usage and pod overhead).
// +optional
Memory *MemoryStats `json:"memory,omitempty"`
// Stats pertaining to network resources.
// +optional
Network *NetworkStats `json:"network,omitempty"`
// Stats pertaining to volume usage of filesystem resources.
// VolumeStats.UsedBytes is the number of bytes used by the Volume
// +optional
// +patchMergeKey=name
// +patchStrategy=merge
VolumeStats []VolumeStats `json:"volume,omitempty" patchStrategy:"merge" patchMergeKey:"name"`
// EphemeralStorage reports the total filesystem usage for the containers and emptyDir-backed volumes in the measured Pod.
// +optional
EphemeralStorage *FsStats `json:"ephemeral-storage,omitempty"`
// ProcessStats pertaining to processes.
// +optional
ProcessStats *ProcessStats `json:"process_stats,omitempty"`
}
// ContainerStats holds container-level unprocessed sample stats.
type ContainerStats struct {
// Reference to the measured container.
Name string `json:"name"`
// The time at which data collection for this container was (re)started.
StartTime metav1.Time `json:"startTime"`
// Stats pertaining to CPU resources.
// +optional
CPU *CPUStats `json:"cpu,omitempty"`
// Stats pertaining to memory (RAM) resources.
// +optional
Memory *MemoryStats `json:"memory,omitempty"`
// Metrics for Accelerators. Each Accelerator corresponds to one element in the array.
Accelerators []AcceleratorStats `json:"accelerators,omitempty"`
// Stats pertaining to container rootfs usage of filesystem resources.
// Rootfs.UsedBytes is the number of bytes used for the container write layer.
// +optional
Rootfs *FsStats `json:"rootfs,omitempty"`
// Stats pertaining to container logs usage of filesystem resources.
// Logs.UsedBytes is the number of bytes used for the container logs.
// +optional
Logs *FsStats `json:"logs,omitempty"`
// User defined metrics that are exposed by containers in the pod. Typically, we expect only one container in the pod to be exposing user defined metrics. In the event of multiple containers exposing metrics, they will be combined here.
// +patchMergeKey=name
// +patchStrategy=merge
UserDefinedMetrics []UserDefinedMetric `json:"userDefinedMetrics,omitempty" patchStrategy:"merge" patchMergeKey:"name"`
}
// PodReference contains enough information to locate the referenced pod.
type PodReference struct {
Name string `json:"name"`
Namespace string `json:"namespace"`
UID string `json:"uid"`
}
// InterfaceStats contains resource value data about interface.
type InterfaceStats struct {
// The name of the interface
Name string `json:"name"`
// Cumulative count of bytes received.
// +optional
RxBytes *uint64 `json:"rxBytes,omitempty"`
// Cumulative count of receive errors encountered.
// +optional
RxErrors *uint64 `json:"rxErrors,omitempty"`
// Cumulative count of bytes transmitted.
// +optional
TxBytes *uint64 `json:"txBytes,omitempty"`
// Cumulative count of transmit errors encountered.
// +optional
TxErrors *uint64 `json:"txErrors,omitempty"`
}
// NetworkStats contains data about network resources.
type NetworkStats struct {
// The time at which these stats were updated.
Time metav1.Time `json:"time"`
// Stats for the default interface, if found
InterfaceStats `json:",inline"`
Interfaces []InterfaceStats `json:"interfaces,omitempty"`
}
// CPUStats contains data about CPU usage.
type CPUStats struct {
// The time at which these stats were updated.
Time metav1.Time `json:"time"`
// Total CPU usage (sum of all cores) averaged over the sample window.
// The "core" unit can be interpreted as CPU core-nanoseconds per second.
// +optional
UsageNanoCores *uint64 `json:"usageNanoCores,omitempty"`
// Cumulative CPU usage (sum of all cores) since object creation.
// +optional
UsageCoreNanoSeconds *uint64 `json:"usageCoreNanoSeconds,omitempty"`
}
// MemoryStats contains data about memory usage.
type MemoryStats struct {
// The time at which these stats were updated.
Time metav1.Time `json:"time"`
// Available memory for use. This is defined as the memory limit - workingSetBytes.
// If memory limit is undefined, the available bytes is omitted.
// +optional
AvailableBytes *uint64 `json:"availableBytes,omitempty"`
// Total memory in use. This includes all memory regardless of when it was accessed.
// +optional
UsageBytes *uint64 `json:"usageBytes,omitempty"`
// The amount of working set memory. This includes recently accessed memory,
// dirty memory, and kernel memory. WorkingSetBytes is <= UsageBytes
// +optional
WorkingSetBytes *uint64 `json:"workingSetBytes,omitempty"`
// The amount of anonymous and swap cache memory (includes transparent
// hugepages).
// +optional
RSSBytes *uint64 `json:"rssBytes,omitempty"`
// Cumulative number of minor page faults.
// +optional
PageFaults *uint64 `json:"pageFaults,omitempty"`
// Cumulative number of major page faults.
// +optional
MajorPageFaults *uint64 `json:"majorPageFaults,omitempty"`
}
// AcceleratorStats contains stats for accelerators attached to the container.
type AcceleratorStats struct {
// Make of the accelerator (nvidia, amd, google etc.)
Make string `json:"make"`
// Model of the accelerator (tesla-p100, tesla-k80 etc.)
Model string `json:"model"`
// ID of the accelerator.
ID string `json:"id"`
// Total accelerator memory.
// unit: bytes
MemoryTotal uint64 `json:"memoryTotal"`
// Total accelerator memory allocated.
// unit: bytes
MemoryUsed uint64 `json:"memoryUsed"`
// Percent of time over the past sample period (10s) during which
// the accelerator was actively processing.
DutyCycle uint64 `json:"dutyCycle"`
}
// VolumeStats contains data about Volume filesystem usage.
type VolumeStats struct {
// Embedded FsStats
FsStats
// Name is the name given to the Volume
// +optional
Name string `json:"name,omitempty"`
// Reference to the PVC, if one exists
// +optional
PVCRef *PVCReference `json:"pvcRef,omitempty"`
}
// PVCReference contains enough information to describe the referenced PVC.
type PVCReference struct {
Name string `json:"name"`
Namespace string `json:"namespace"`
}
// FsStats contains data about filesystem usage.
type FsStats struct {
// The time at which these stats were updated.
Time metav1.Time `json:"time"`
// AvailableBytes represents the storage space available (bytes) for the filesystem.
// +optional
AvailableBytes *uint64 `json:"availableBytes,omitempty"`
// CapacityBytes represents the total capacity (bytes) of the filesystems underlying storage.
// +optional
CapacityBytes *uint64 `json:"capacityBytes,omitempty"`
// UsedBytes represents the bytes used for a specific task on the filesystem.
// This may differ from the total bytes used on the filesystem and may not equal CapacityBytes - AvailableBytes.
// e.g. For ContainerStats.Rootfs this is the bytes used by the container rootfs on the filesystem.
// +optional
UsedBytes *uint64 `json:"usedBytes,omitempty"`
// InodesFree represents the free inodes in the filesystem.
// +optional
InodesFree *uint64 `json:"inodesFree,omitempty"`
// Inodes represents the total inodes in the filesystem.
// +optional
Inodes *uint64 `json:"inodes,omitempty"`
// InodesUsed represents the inodes used by the filesystem
// This may not equal Inodes - InodesFree because this filesystem may share inodes with other "filesystems"
// e.g. For ContainerStats.Rootfs, this is the inodes used only by that container, and does not count inodes used by other containers.
InodesUsed *uint64 `json:"inodesUsed,omitempty"`
}
// UserDefinedMetricType defines how the metric should be interpreted by the user.
type UserDefinedMetricType string
const (
// MetricGauge is an instantaneous value. May increase or decrease.
MetricGauge UserDefinedMetricType = "gauge"
// MetricCumulative is a counter-like value that is only expected to increase.
MetricCumulative UserDefinedMetricType = "cumulative"
// MetricDelta is a rate over a time period.
MetricDelta UserDefinedMetricType = "delta"
)
// UserDefinedMetricDescriptor contains metadata that describes a user defined metric.
type UserDefinedMetricDescriptor struct {
// The name of the metric.
Name string `json:"name"`
// Type of the metric.
Type UserDefinedMetricType `json:"type"`
// Display Units for the stats.
Units string `json:"units"`
// Metadata labels associated with this metric.
// +optional
Labels map[string]string `json:"labels,omitempty"`
}
// UserDefinedMetric represents a metric defined and generated by users.
type UserDefinedMetric struct {
UserDefinedMetricDescriptor `json:",inline"`
// The time at which these stats were updated.
Time metav1.Time `json:"time"`
// Value of the metric. Float64s have 53 bit precision.
// We do not foresee any metrics exceeding that value.
Value float64 `json:"value"`
}
// ProcessStats are stats pertaining to processes.
type ProcessStats struct {
// Number of processes
// +optional
ProcessCount *uint64 `json:"process_count,omitempty"`
}
const (
KubernetesPodNameLabel = "io.kubernetes.pod.name"
KubernetesPodNamespaceLabel = "io.kubernetes.pod.namespace"
KubernetesPodUIDLabel = "io.kubernetes.pod.uid"
KubernetesContainerNameLabel = "io.kubernetes.container.name"
)
func GetContainerName(labels map[string]string) string {
return labels[KubernetesContainerNameLabel]
}
func GetPodName(labels map[string]string) string {
return labels[KubernetesPodNameLabel]
}
func GetPodUID(labels map[string]string) string {
return labels[KubernetesPodUIDLabel]
}
func GetPodNamespace(labels map[string]string) string {
return labels[KubernetesPodNamespaceLabel]
}