mirror of
https://github.com/Tencent/WeKnora.git
synced 2026-10-02 05:54:33 +08:00
feat(plugin): report the memory each plugin uses
The platform admin could not see what plugins cost: how much memory a
plugin's processes take on each node, or whether an idle one is stopped.
The host measures the resident memory of each plugin's process group,
which holds its children and, on Linux, the network sandbox's relay: from
/proc on Linux, from one ps run elsewhere (not on Windows). One sample of
all plugins is reused for 10 seconds; a plugin that just started is
measured again after a second. The host reports it through a new
reconcile.UsageReporter, with whether the plugin is stopped while idle,
and each node publishes it with its status; a single node reads its own.
The admin list sums it up per plugin in a memory field (bytes over the
measured instances, and how many are idle), shown in a new memory column
("idle" when all instances are stopped, "-" for remote and kubernetes
plugins). The detail drawer shows each node's memory or an idle tag.
Checked against ps on the example RSS plugin: 11812864 bytes reported
for 11536 KB.
This commit is contained in:
@@ -120,6 +120,10 @@ export interface PluginInstance {
|
||||
upgradeVersion?: string
|
||||
upgradeState?: UpgradeState
|
||||
upgradeError?: string
|
||||
/** Resident memory of the instance's processes; absent when not measured. */
|
||||
memoryBytes?: number
|
||||
/** Stopped for going without calls; the next call starts it. */
|
||||
idle?: boolean
|
||||
}
|
||||
|
||||
/** Where an upgrade an instance has not loaded stands. */
|
||||
|
||||
@@ -58,6 +58,17 @@ export interface PluginNodeStatus {
|
||||
upgradeVersion?: string
|
||||
upgradeState?: UpgradeState
|
||||
upgradeError?: string
|
||||
memoryBytes?: number
|
||||
idle?: boolean
|
||||
}
|
||||
|
||||
/** The resident memory of a plugin's instances, summed up. */
|
||||
export interface PluginMemory {
|
||||
/** Total over the measured instances. */
|
||||
bytes: number
|
||||
measured: number
|
||||
/** Instances stopped for going without calls. */
|
||||
idle: number
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -91,6 +102,8 @@ export interface InstalledPlugin {
|
||||
egress?: EgressMode
|
||||
/** Set while instances have not loaded the active version. */
|
||||
upgrade?: PluginUpgrade
|
||||
/** Absent when no instance is measured or idle (remote, kubernetes). */
|
||||
memory?: PluginMemory
|
||||
/** Set for a workspace's own plugin: the workspace that registered it. */
|
||||
owner_tenant_id?: number
|
||||
/** The workspaces the plugin is limited to; absent when every workspace sees it. */
|
||||
|
||||
@@ -2748,7 +2748,8 @@ export default {
|
||||
runtime: 'Runtime',
|
||||
trust: 'Trust',
|
||||
audience: 'Visible to',
|
||||
state: 'Status'
|
||||
state: 'Status',
|
||||
memory: 'Memory'
|
||||
},
|
||||
menu: {
|
||||
detail: 'View details',
|
||||
@@ -2811,6 +2812,12 @@ export default {
|
||||
pendingNode: 'Upgrading to v{version}',
|
||||
failedNode: 'Upgrade to v{version} failed'
|
||||
},
|
||||
memory: {
|
||||
idle: 'Idle',
|
||||
idleHint: '{count} instance(s) stopped for going without calls; the next call starts them',
|
||||
total: 'Total of {count} instance(s)',
|
||||
node: 'Memory {size}'
|
||||
},
|
||||
source: {
|
||||
upload: 'Upload',
|
||||
url: 'URL'
|
||||
|
||||
@@ -2748,7 +2748,8 @@ export default {
|
||||
runtime: '実行方式',
|
||||
trust: '信頼',
|
||||
audience: '公開範囲',
|
||||
state: '状態'
|
||||
state: '状態',
|
||||
memory: 'メモリ'
|
||||
},
|
||||
menu: {
|
||||
detail: '詳細を表示',
|
||||
@@ -2811,6 +2812,12 @@ export default {
|
||||
pendingNode: 'v{version} へアップグレード中',
|
||||
failedNode: 'v{version} へのアップグレードに失敗'
|
||||
},
|
||||
memory: {
|
||||
idle: 'アイドル',
|
||||
idleHint: '{count} 個のインスタンスは呼び出しがないため停止中。次の呼び出しで起動します',
|
||||
total: '{count} 個のインスタンスの合計',
|
||||
node: 'メモリ {size}'
|
||||
},
|
||||
source: {
|
||||
upload: 'アップロード',
|
||||
url: 'URL'
|
||||
|
||||
@@ -5274,7 +5274,8 @@ export default {
|
||||
runtime: '실행 방식',
|
||||
trust: '신뢰',
|
||||
audience: '공개 범위',
|
||||
state: '상태'
|
||||
state: '상태',
|
||||
memory: '메모리'
|
||||
},
|
||||
menu: {
|
||||
detail: '상세 보기',
|
||||
@@ -5337,6 +5338,12 @@ export default {
|
||||
pendingNode: 'v{version}(으)로 업그레이드 중',
|
||||
failedNode: 'v{version}(으)로 업그레이드 실패'
|
||||
},
|
||||
memory: {
|
||||
idle: '유휴',
|
||||
idleHint: '인스턴스 {count}개가 호출이 없어 중지됨. 다음 호출 시 시작됩니다',
|
||||
total: '인스턴스 {count}개 합계',
|
||||
node: '메모리 {size}'
|
||||
},
|
||||
source: {
|
||||
upload: '업로드',
|
||||
url: 'URL'
|
||||
|
||||
@@ -5274,7 +5274,8 @@ export default {
|
||||
runtime: 'Запуск',
|
||||
trust: 'Доверие',
|
||||
audience: 'Доступен',
|
||||
state: 'Статус'
|
||||
state: 'Статус',
|
||||
memory: 'Память'
|
||||
},
|
||||
menu: {
|
||||
detail: 'Подробнее',
|
||||
@@ -5337,6 +5338,12 @@ export default {
|
||||
pendingNode: 'Обновление до v{version}',
|
||||
failedNode: 'Не удалось обновить до v{version}'
|
||||
},
|
||||
memory: {
|
||||
idle: 'Простаивает',
|
||||
idleHint: 'Экземпляров остановлено без вызовов: {count}; следующий вызов их запустит',
|
||||
total: 'Всего по экземплярам: {count}',
|
||||
node: 'Память {size}'
|
||||
},
|
||||
source: {
|
||||
upload: 'Загрузка',
|
||||
url: 'URL'
|
||||
|
||||
@@ -5276,7 +5276,8 @@ export default {
|
||||
runtime: '运行方式',
|
||||
trust: '信任',
|
||||
audience: '可见范围',
|
||||
state: '状态'
|
||||
state: '状态',
|
||||
memory: '内存'
|
||||
},
|
||||
menu: {
|
||||
detail: '查看详情',
|
||||
@@ -5339,6 +5340,12 @@ export default {
|
||||
pendingNode: '正在升级到 v{version}',
|
||||
failedNode: '升级到 v{version} 失败'
|
||||
},
|
||||
memory: {
|
||||
idle: '空闲',
|
||||
idleHint: '{count} 个实例因没有调用已停止进程,下次调用时启动',
|
||||
total: '{count} 个实例合计',
|
||||
node: '内存 {size}'
|
||||
},
|
||||
source: {
|
||||
upload: '上传',
|
||||
url: 'URL'
|
||||
|
||||
@@ -50,13 +50,14 @@
|
||||
<th>{{ t('pluginAdmin.columns.runtime') }}</th>
|
||||
<th>{{ t('pluginAdmin.columns.trust') }}</th>
|
||||
<th>{{ t('pluginAdmin.columns.audience') }}</th>
|
||||
<th>{{ t('pluginAdmin.columns.memory') }}</th>
|
||||
<th>{{ t('pluginAdmin.columns.state') }}</th>
|
||||
<th class="plugin-table__menu-col" />
|
||||
</tr>
|
||||
</thead>
|
||||
<tbody>
|
||||
<tr v-if="visible.length === 0">
|
||||
<td colspan="7" class="plugin-table__empty">{{ t('pluginAdmin.noMatch') }}</td>
|
||||
<td colspan="8" class="plugin-table__empty">{{ t('pluginAdmin.noMatch') }}</td>
|
||||
</tr>
|
||||
<tr
|
||||
v-for="p in visible"
|
||||
@@ -83,6 +84,21 @@
|
||||
<span class="trust" :class="`trust--${activeTrust(p)}`">{{ t(`pluginAdmin.trust.${activeTrust(p)}`) }}</span>
|
||||
</td>
|
||||
<td>{{ audienceLabel(p) }}</td>
|
||||
<td class="plugin-row__num">
|
||||
<template v-for="m in [memoryCell(p)]" :key="m.kind">
|
||||
<t-tooltip v-if="m.kind === 'bytes'">
|
||||
<template #content>
|
||||
<div>{{ t('pluginAdmin.memory.total', { count: m.measured }) }}</div>
|
||||
<div v-if="m.idle">{{ t('pluginAdmin.memory.idleHint', { count: m.idle }) }}</div>
|
||||
</template>
|
||||
<span>{{ formatBytes(m.bytes) }}</span>
|
||||
</t-tooltip>
|
||||
<t-tooltip v-else-if="m.kind === 'idle'" :content="t('pluginAdmin.memory.idleHint', { count: m.idle })">
|
||||
<span class="memory-idle">{{ t('pluginAdmin.memory.idle') }}</span>
|
||||
</t-tooltip>
|
||||
<span v-else class="memory-none">—</span>
|
||||
</template>
|
||||
</td>
|
||||
<td>
|
||||
<div class="plugin-row__state">
|
||||
<t-tooltip :content="p.node?.error" :disabled="!p.node?.error">
|
||||
@@ -154,6 +170,8 @@ import {
|
||||
activeTrust,
|
||||
audienceSummary,
|
||||
egressUnenforced,
|
||||
formatBytes,
|
||||
memoryCell,
|
||||
filterInstalled,
|
||||
installedState,
|
||||
matchesAdminFilter,
|
||||
@@ -388,6 +406,11 @@ onMounted(load)
|
||||
font-variant-numeric: tabular-nums;
|
||||
}
|
||||
|
||||
.memory-idle,
|
||||
.memory-none {
|
||||
color: var(--td-text-color-placeholder);
|
||||
}
|
||||
|
||||
.plugin-row__state {
|
||||
display: inline-flex;
|
||||
align-items: center;
|
||||
|
||||
@@ -17,6 +17,7 @@ import {
|
||||
trustNote,
|
||||
trustTheme,
|
||||
upgradeTheme,
|
||||
memoryCell,
|
||||
filterMarket,
|
||||
marketAction,
|
||||
audienceOf,
|
||||
@@ -200,3 +201,13 @@ test('upgrade themes', () => {
|
||||
assert.equal(upgradeTheme('failed'), 'warning')
|
||||
assert.equal(upgradeTheme('pending'), 'primary')
|
||||
})
|
||||
|
||||
test('memory column', () => {
|
||||
const on = 'enabled' as const
|
||||
assert.deepEqual(memoryCell({ desired_state: on }), { kind: 'none' })
|
||||
assert.deepEqual(memoryCell({ desired_state: 'disabled', memory: { bytes: 1, measured: 1, idle: 0 } }), { kind: 'none' })
|
||||
assert.deepEqual(memoryCell({ desired_state: on, memory: { bytes: 5, measured: 2, idle: 1 } }), {
|
||||
kind: 'bytes', bytes: 5, measured: 2, idle: 1,
|
||||
})
|
||||
assert.deepEqual(memoryCell({ desired_state: on, memory: { bytes: 0, measured: 0, idle: 3 } }), { kind: 'idle', idle: 3 })
|
||||
})
|
||||
|
||||
@@ -159,6 +159,20 @@ export function installedState(p: InstalledPlugin): InstalledState {
|
||||
}
|
||||
}
|
||||
|
||||
/** What the memory column shows for a plugin. */
|
||||
export type MemoryCell =
|
||||
| { kind: 'bytes'; bytes: number; measured: number; idle: number }
|
||||
| { kind: 'idle'; idle: number }
|
||||
| { kind: 'none' }
|
||||
|
||||
export function memoryCell(p: Pick<InstalledPlugin, 'memory' | 'desired_state'>): MemoryCell {
|
||||
const m = p.memory
|
||||
if (p.desired_state === 'disabled' || !m) return { kind: 'none' }
|
||||
if (m.measured > 0) return { kind: 'bytes', bytes: m.bytes, measured: m.measured, idle: m.idle }
|
||||
if (m.idle > 0) return { kind: 'idle', idle: m.idle }
|
||||
return { kind: 'none' }
|
||||
}
|
||||
|
||||
/** The tag theme of an unfinished upgrade: a failed one warns. */
|
||||
export function upgradeTheme(state: UpgradeState | undefined): 'warning' | 'primary' {
|
||||
return state === 'failed' ? 'warning' : 'primary'
|
||||
|
||||
@@ -71,6 +71,12 @@
|
||||
</t-tag>
|
||||
<code>{{ n.node }}</code>
|
||||
<span class="line-list__muted">v{{ n.version }}</span>
|
||||
<t-tooltip v-if="n.idle" :content="t('pluginAdmin.memory.idleHint', { count: 1 })">
|
||||
<t-tag size="small" variant="outline">{{ t('pluginAdmin.memory.idle') }}</t-tag>
|
||||
</t-tooltip>
|
||||
<span v-else-if="n.memoryBytes" class="line-list__muted">
|
||||
{{ t('pluginAdmin.memory.node', { size: formatBytes(n.memoryBytes) }) }}
|
||||
</span>
|
||||
<t-tooltip v-if="n.egress" :content="t(`pluginAdmin.egress.hint.${n.egress}`)">
|
||||
<t-tag size="small" variant="outline" :theme="egressTheme(n.egress)">
|
||||
{{ t(`pluginAdmin.egress.mode.${n.egress}`) }}
|
||||
|
||||
@@ -276,6 +276,18 @@ type InstalledPluginDTO struct {
|
||||
// Upgrade is set while instances have not loaded the active version:
|
||||
// still starting it, or failed to. They keep running the previous one.
|
||||
Upgrade *PluginUpgradeDTO `json:"upgrade,omitempty"`
|
||||
// Memory sums up what the plugin's processes use; nil when no instance
|
||||
// is measured or idle (remote and kubernetes plugins).
|
||||
Memory *PluginMemoryDTO `json:"memory,omitempty"`
|
||||
}
|
||||
|
||||
// PluginMemoryDTO sums up the resident memory of a plugin's instances.
|
||||
type PluginMemoryDTO struct {
|
||||
// Bytes is the total over the measured instances.
|
||||
Bytes int64 `json:"bytes"`
|
||||
Measured int `json:"measured"`
|
||||
// Idle counts instances stopped for going without calls.
|
||||
Idle int `json:"idle"`
|
||||
}
|
||||
|
||||
// PluginUpgradeDTO sums up the instances that have not loaded the active
|
||||
@@ -292,8 +304,8 @@ type PluginUpgradeDTO struct {
|
||||
}
|
||||
|
||||
// withEgress adds what the plugin's instances report: how its outbound
|
||||
// traffic is controlled across the nodes and plugin hosts running it, and
|
||||
// an upgrade they have not finished.
|
||||
// traffic is controlled across the nodes and plugin hosts running it, an
|
||||
// upgrade they have not finished, and the memory they use.
|
||||
func (h *PluginAdminHandler) withEgress(ctx context.Context, v *install.View) InstalledPluginDTO {
|
||||
out := InstalledPluginDTO{View: v}
|
||||
if h.drivers == nil || v.DesiredState != types.PluginStateEnabled ||
|
||||
@@ -311,9 +323,29 @@ func (h *PluginAdminHandler) withEgress(ctx context.Context, v *install.View) In
|
||||
}
|
||||
out.Egress = weakestEgress(instances)
|
||||
out.Upgrade = upgradeOf(instances)
|
||||
out.Memory = memoryOf(instances)
|
||||
return out
|
||||
}
|
||||
|
||||
// memoryOf sums up the instances' memory; nil when none is measured or
|
||||
// idle.
|
||||
func memoryOf(instances []driver.InstanceStatus) *PluginMemoryDTO {
|
||||
var out PluginMemoryDTO
|
||||
for _, in := range instances {
|
||||
switch {
|
||||
case in.Idle:
|
||||
out.Idle++
|
||||
case in.MemoryBytes > 0:
|
||||
out.Measured++
|
||||
out.Bytes += in.MemoryBytes
|
||||
}
|
||||
}
|
||||
if out.Measured == 0 && out.Idle == 0 {
|
||||
return nil
|
||||
}
|
||||
return &out
|
||||
}
|
||||
|
||||
// upgradeOf sums up instances with an upgrade they have not loaded; nil
|
||||
// when there are none.
|
||||
func upgradeOf(instances []driver.InstanceStatus) *PluginUpgradeDTO {
|
||||
|
||||
@@ -310,3 +310,15 @@ func TestUpgradeOf(t *testing.T) {
|
||||
t.Fatalf("upgrade = %+v", u)
|
||||
}
|
||||
}
|
||||
|
||||
func TestMemoryOf(t *testing.T) {
|
||||
if memoryOf([]driver.InstanceStatus{{Node: "remote"}}) != nil {
|
||||
t.Fatal("instances without a measurement sum up to nothing")
|
||||
}
|
||||
got := memoryOf([]driver.InstanceStatus{
|
||||
{MemoryBytes: 10 << 20}, {MemoryBytes: 6 << 20}, {Idle: true}, {},
|
||||
})
|
||||
if got == nil || *got != (PluginMemoryDTO{Bytes: 16 << 20, Measured: 2, Idle: 1}) {
|
||||
t.Fatalf("memory = %+v", got)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -86,6 +86,11 @@ type InstanceStatus struct {
|
||||
UpgradeVersion string `json:"upgradeVersion,omitempty"`
|
||||
UpgradeState string `json:"upgradeState,omitempty"`
|
||||
UpgradeError string `json:"upgradeError,omitempty"`
|
||||
// MemoryBytes is the resident memory of the instance's processes; 0
|
||||
// when not measured. Idle: stopped for going without calls, it starts
|
||||
// on the next one.
|
||||
MemoryBytes int64 `json:"memoryBytes,omitempty"`
|
||||
Idle bool `json:"idle,omitempty"`
|
||||
}
|
||||
|
||||
// Driver runs plugins of one runtime type.
|
||||
|
||||
@@ -64,6 +64,8 @@ type Manager struct {
|
||||
idleTimeout time.Duration
|
||||
reaperOnce sync.Once
|
||||
reaperStop chan struct{}
|
||||
|
||||
memory *memorySampler
|
||||
}
|
||||
|
||||
// handoverWindow is how long a standalone host keeps serving the previous
|
||||
@@ -74,7 +76,10 @@ var handoverWindow = reconcile.DefaultInterval + 15*time.Second
|
||||
|
||||
// NewManager creates an empty host that runs every kind this machine can.
|
||||
func NewManager() *Manager {
|
||||
m := &Manager{procs: map[string]*process{}, idleTimeout: IdleTimeoutFromEnv(), reaperStop: make(chan struct{})}
|
||||
m := &Manager{
|
||||
procs: map[string]*process{}, idleTimeout: IdleTimeoutFromEnv(), reaperStop: make(chan struct{}),
|
||||
memory: newMemorySampler(),
|
||||
}
|
||||
m.SetKinds(AvailableKinds())
|
||||
return m
|
||||
}
|
||||
@@ -84,7 +89,7 @@ func NewManager() *Manager {
|
||||
func NewStandaloneManager(kinds []string) *Manager {
|
||||
m := &Manager{
|
||||
procs: map[string]*process{}, standalone: true, handingOver: map[string]*process{},
|
||||
idleTimeout: IdleTimeoutFromEnv(), reaperStop: make(chan struct{}),
|
||||
idleTimeout: IdleTimeoutFromEnv(), reaperStop: make(chan struct{}), memory: newMemorySampler(),
|
||||
}
|
||||
m.SetKinds(kinds)
|
||||
return m
|
||||
|
||||
@@ -0,0 +1,87 @@
|
||||
package host
|
||||
|
||||
import (
|
||||
"context"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
"github.com/Tencent/WeKnora/internal/logger"
|
||||
"github.com/Tencent/WeKnora/internal/plugin/reconcile"
|
||||
)
|
||||
|
||||
// memorySampleAge is how long one memory sample of all plugins is reused:
|
||||
// status reports come in bursts, and ps is a process on some systems. A
|
||||
// group missing from a sample (a plugin that just started) is measured
|
||||
// again once the sample is memoryMissAge old.
|
||||
const (
|
||||
memorySampleAge = 10 * time.Second
|
||||
memoryMissAge = time.Second
|
||||
)
|
||||
|
||||
// memorySampler measures the plugins' process groups together and keeps
|
||||
// the result for memorySampleAge.
|
||||
type memorySampler struct {
|
||||
mu sync.Mutex
|
||||
at time.Time
|
||||
byGroup map[int]int64
|
||||
warned bool
|
||||
measure func(map[int]bool) (map[int]int64, error)
|
||||
now func() time.Time
|
||||
}
|
||||
|
||||
func newMemorySampler() *memorySampler {
|
||||
return &memorySampler{measure: groupMemory, now: time.Now}
|
||||
}
|
||||
|
||||
// group returns a process group's resident memory, measuring the groups of
|
||||
// every running plugin when the last sample is too old.
|
||||
func (s *memorySampler) group(pgid int, all func() map[int]bool) int64 {
|
||||
s.mu.Lock()
|
||||
defer s.mu.Unlock()
|
||||
age := s.now().Sub(s.at)
|
||||
_, seen := s.byGroup[pgid]
|
||||
if s.byGroup == nil || age >= memorySampleAge || (!seen && age >= memoryMissAge) {
|
||||
sums, err := s.measure(all())
|
||||
if err != nil && !s.warned {
|
||||
s.warned = true
|
||||
logger.Warnf(context.Background(), "[plugin] cannot measure plugin memory: %v", err)
|
||||
}
|
||||
s.byGroup, s.at = sums, s.now()
|
||||
}
|
||||
return s.byGroup[pgid]
|
||||
}
|
||||
|
||||
// Usage implements reconcile.UsageReporter: the resident memory of a
|
||||
// plugin's processes on this host, or that it is stopped while idle.
|
||||
func (m *Manager) Usage(pluginID string) (reconcile.Usage, bool) {
|
||||
m.mu.Lock()
|
||||
p := m.procs[pluginID]
|
||||
m.mu.Unlock()
|
||||
if p == nil {
|
||||
return reconcile.Usage{}, false
|
||||
}
|
||||
p.mu.RLock()
|
||||
idle := p.state == StateIdle
|
||||
p.mu.RUnlock()
|
||||
if idle {
|
||||
return reconcile.Usage{Idle: true}, true
|
||||
}
|
||||
pgid := int(p.group.Load())
|
||||
if pgid == 0 {
|
||||
return reconcile.Usage{}, true
|
||||
}
|
||||
return reconcile.Usage{MemoryBytes: m.memory.group(pgid, m.groups)}, true
|
||||
}
|
||||
|
||||
// groups are the process groups of the plugins running on this host.
|
||||
func (m *Manager) groups() map[int]bool {
|
||||
m.mu.Lock()
|
||||
defer m.mu.Unlock()
|
||||
out := make(map[int]bool, len(m.procs))
|
||||
for _, p := range m.procs {
|
||||
if g := p.group.Load(); g != 0 {
|
||||
out[int(g)] = true
|
||||
}
|
||||
}
|
||||
return out
|
||||
}
|
||||
@@ -0,0 +1,52 @@
|
||||
//go:build linux
|
||||
|
||||
package host
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"os"
|
||||
"strconv"
|
||||
)
|
||||
|
||||
// groupMemory sums the resident memory of the processes in each of the
|
||||
// given process groups, read from /proc.
|
||||
func groupMemory(groups map[int]bool) (map[int]int64, error) {
|
||||
entries, err := os.ReadDir("/proc")
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
page := int64(os.Getpagesize())
|
||||
out := make(map[int]int64, len(groups))
|
||||
for _, e := range entries {
|
||||
if _, err := strconv.Atoi(e.Name()); err != nil {
|
||||
continue
|
||||
}
|
||||
stat, err := os.ReadFile("/proc/" + e.Name() + "/stat")
|
||||
if err != nil {
|
||||
continue // exited meanwhile
|
||||
}
|
||||
pgrp, rss, ok := parseStat(stat)
|
||||
if ok && groups[pgrp] {
|
||||
out[pgrp] += rss * page
|
||||
}
|
||||
}
|
||||
return out, nil
|
||||
}
|
||||
|
||||
// parseStat reads the process group and resident pages from a
|
||||
// /proc/<pid>/stat line. The command name may hold spaces and parentheses,
|
||||
// so fields are counted after its last ')'.
|
||||
func parseStat(stat []byte) (pgrp int, rssPages int64, ok bool) {
|
||||
i := bytes.LastIndexByte(stat, ')')
|
||||
if i < 0 {
|
||||
return 0, 0, false
|
||||
}
|
||||
// After the name: state ppid pgrp ... with rss the 22nd of them.
|
||||
f := bytes.Fields(stat[i+1:])
|
||||
if len(f) < 22 {
|
||||
return 0, 0, false
|
||||
}
|
||||
pgrp, err1 := strconv.Atoi(string(f[2]))
|
||||
rss, err2 := strconv.ParseInt(string(f[21]), 10, 64)
|
||||
return pgrp, rss, err1 == nil && err2 == nil
|
||||
}
|
||||
@@ -0,0 +1,16 @@
|
||||
package host
|
||||
|
||||
import "testing"
|
||||
|
||||
func TestParseStat(t *testing.T) {
|
||||
// A name with spaces and a parenthesis; pgrp 4242, rss 1536 pages.
|
||||
line := "4250 (weird (name) x) S 4242 4242 4242 0 -1 4194560 100 0 0 0 1 2 0 0 20 0 3 0 1000 " +
|
||||
"123456789 1536 18446744073709551615 1 1 0 0 0 0 0 0 0 0 0 0 17 3 0 0 0 0 0"
|
||||
pgrp, rss, ok := parseStat([]byte(line))
|
||||
if !ok || pgrp != 4242 || rss != 1536 {
|
||||
t.Fatalf("parseStat = %d, %d, %v", pgrp, rss, ok)
|
||||
}
|
||||
if _, _, ok := parseStat([]byte("garbage")); ok {
|
||||
t.Fatal("garbage parsed")
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,77 @@
|
||||
package host
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"runtime"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/Tencent/WeKnora/internal/plugin/reconcile"
|
||||
)
|
||||
|
||||
// One sample serves every plugin until it is memorySampleAge old.
|
||||
func TestMemorySampleIsShared(t *testing.T) {
|
||||
now := time.Unix(1000, 0)
|
||||
calls := 0
|
||||
s := &memorySampler{
|
||||
now: func() time.Time { return now },
|
||||
measure: func(map[int]bool) (map[int]int64, error) {
|
||||
calls++
|
||||
return map[int]int64{1: 10 << 20, 2: 20 << 20}, nil
|
||||
},
|
||||
}
|
||||
all := func() map[int]bool { return map[int]bool{1: true, 2: true} }
|
||||
if s.group(1, all) != 10<<20 || s.group(2, all) != 20<<20 || calls != 1 {
|
||||
t.Fatalf("calls = %d", calls)
|
||||
}
|
||||
now = now.Add(memorySampleAge)
|
||||
s.group(1, all)
|
||||
if calls != 2 {
|
||||
t.Fatalf("an old sample was reused: calls = %d", calls)
|
||||
}
|
||||
// A group the sample lacks (a plugin that just started) is measured
|
||||
// again, but not more than once a memoryMissAge.
|
||||
s.group(3, all)
|
||||
if calls != 2 {
|
||||
t.Fatalf("a fresh sample was measured again for a new group: calls = %d", calls)
|
||||
}
|
||||
now = now.Add(memoryMissAge)
|
||||
s.group(3, all)
|
||||
if calls != 3 {
|
||||
t.Fatalf("a missing group was not measured again: calls = %d", calls)
|
||||
}
|
||||
s.measure = func(map[int]bool) (map[int]int64, error) { return nil, errors.New("no ps") }
|
||||
now = now.Add(memorySampleAge)
|
||||
if got := s.group(1, all); got != 0 {
|
||||
t.Fatalf("unmeasured memory = %d", got)
|
||||
}
|
||||
}
|
||||
|
||||
// A running plugin reports the memory of its process; an idle one says so.
|
||||
func TestUsageReportsMemoryAndIdle(t *testing.T) {
|
||||
if runtime.GOOS == "windows" {
|
||||
t.Skip("plugin memory is not measured on Windows")
|
||||
}
|
||||
fastTimings(t)
|
||||
ctx := context.Background()
|
||||
m := NewManager()
|
||||
m.SetIdleTimeout(300 * time.Millisecond)
|
||||
defer m.Close()
|
||||
if _, ok := m.Usage("acme.echo"); ok {
|
||||
t.Fatal("usage of a plugin this host does not run")
|
||||
}
|
||||
if err := reconcile.Activate(ctx, m, install(t, "1.0.0", "")); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
u, ok := m.Usage("acme.echo")
|
||||
// The echo plugin is a Go binary: a few MB at least.
|
||||
if !ok || u.Idle || u.MemoryBytes < 1<<20 {
|
||||
t.Fatalf("usage = %+v, %v", u, ok)
|
||||
}
|
||||
var _ reconcile.UsageReporter = m
|
||||
waitFor(t, "the idle plugin to stop", func() bool { return state(m, "acme.echo") == StateIdle })
|
||||
if u, _ := m.Usage("acme.echo"); !u.Idle || u.MemoryBytes != 0 {
|
||||
t.Fatalf("idle usage = %+v", u)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,34 @@
|
||||
//go:build unix && !linux
|
||||
|
||||
package host
|
||||
|
||||
import (
|
||||
"bufio"
|
||||
"bytes"
|
||||
"os/exec"
|
||||
"strconv"
|
||||
"strings"
|
||||
)
|
||||
|
||||
// groupMemory sums the resident memory of the processes in each of the
|
||||
// given process groups, from one run of ps.
|
||||
func groupMemory(groups map[int]bool) (map[int]int64, error) {
|
||||
out, err := exec.Command("ps", "-A", "-o", "pgid=,rss=").Output()
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
sums := make(map[int]int64, len(groups))
|
||||
sc := bufio.NewScanner(bytes.NewReader(out))
|
||||
for sc.Scan() {
|
||||
f := strings.Fields(sc.Text())
|
||||
if len(f) != 2 {
|
||||
continue
|
||||
}
|
||||
pgid, err1 := strconv.Atoi(f[0])
|
||||
kb, err2 := strconv.ParseInt(f[1], 10, 64)
|
||||
if err1 == nil && err2 == nil && groups[pgid] {
|
||||
sums[pgid] += kb * 1024
|
||||
}
|
||||
}
|
||||
return sums, sc.Err()
|
||||
}
|
||||
@@ -0,0 +1,8 @@
|
||||
package host
|
||||
|
||||
import "errors"
|
||||
|
||||
// groupMemory is not measured on Windows.
|
||||
func groupMemory(map[int]bool) (map[int]int64, error) {
|
||||
return nil, errors.New("plugin memory is not measured on Windows")
|
||||
}
|
||||
@@ -16,6 +16,7 @@ import (
|
||||
"runtime"
|
||||
"strings"
|
||||
"sync"
|
||||
"sync/atomic"
|
||||
"time"
|
||||
|
||||
"github.com/Tencent/WeKnora/internal/logger"
|
||||
@@ -147,6 +148,9 @@ type process struct {
|
||||
changed chan struct{}
|
||||
|
||||
activity *activity
|
||||
// group is the process group of the running child (its pid: children
|
||||
// start their own group), 0 while none runs.
|
||||
group atomic.Int64
|
||||
// park asks the supervisor to stop the idle child; wake to start it.
|
||||
park chan struct{}
|
||||
wake chan struct{}
|
||||
@@ -202,6 +206,7 @@ func startProcess(sp spec, onState func(*process, State, error)) (*process, erro
|
||||
return nil, err
|
||||
}
|
||||
// Ready before Activate returns: callers route to the plugin right away.
|
||||
p.group.Store(int64(first.cmd.Process.Pid))
|
||||
p.setClient(first.client)
|
||||
p.setState(StateReady, nil)
|
||||
go p.supervise(ctx, entry, first)
|
||||
@@ -575,9 +580,11 @@ func (p *process) supervise(ctx context.Context, entry string, cur *launched) {
|
||||
defer close(p.done)
|
||||
backoff := restartBackoffFloor
|
||||
for {
|
||||
p.group.Store(int64(cur.cmd.Process.Pid))
|
||||
p.setClient(cur.client)
|
||||
p.setState(StateReady, nil)
|
||||
exitErr := p.watch(ctx, cur)
|
||||
p.group.Store(0)
|
||||
p.setClient(nil)
|
||||
if ctx.Err() != nil {
|
||||
p.setState(StateStopped, nil)
|
||||
|
||||
@@ -56,6 +56,10 @@ func (r *Reconciler) publishStatuses(ctx context.Context) {
|
||||
snapshot[id] = s
|
||||
}
|
||||
r.statusMu.RUnlock()
|
||||
for id, s := range snapshot {
|
||||
s.Usage = r.usage(id)
|
||||
snapshot[id] = s
|
||||
}
|
||||
node, now := r.NodeName(), time.Now()
|
||||
pipe := r.rdb.Pipeline()
|
||||
for id, s := range snapshot {
|
||||
@@ -154,6 +158,7 @@ func (d nodeDriver) Status(ctx context.Context, pluginID string) ([]driver.Insta
|
||||
out = append(out, driver.InstanceStatus{
|
||||
Node: n.Node, Version: n.Version, State: state, Error: n.Error, UpdatedAt: n.UpdatedAt, Egress: egress,
|
||||
UpgradeVersion: n.UpgradeVersion, UpgradeState: n.UpgradeState, UpgradeError: n.UpgradeError,
|
||||
MemoryBytes: n.MemoryBytes, Idle: n.Idle,
|
||||
})
|
||||
}
|
||||
return out, nil
|
||||
|
||||
@@ -130,3 +130,61 @@ func waitFor(t *testing.T, cond func() bool) {
|
||||
time.Sleep(10 * time.Millisecond)
|
||||
}
|
||||
}
|
||||
|
||||
// usageActivator reports fixed usage for every plugin.
|
||||
type usageActivator struct {
|
||||
recorder
|
||||
usage Usage
|
||||
}
|
||||
|
||||
func (a *usageActivator) Usage(string) (Usage, bool) { return a.usage, true }
|
||||
|
||||
// What a plugin's processes use is part of each node's report, read when
|
||||
// the report is made; a single node without Redis shows its own.
|
||||
func TestNodesReportUsage(t *testing.T) {
|
||||
ctx := context.Background()
|
||||
repo, store := plugintest.NewMemRepo(), &plugintest.MemStore{}
|
||||
plugintest.Install(t, repo, store, plugintest.KitPackage(t, "1.0.0"), types.PluginStateEnabled)
|
||||
|
||||
act := &usageActivator{usage: Usage{MemoryBytes: 12 << 20}}
|
||||
single := New(Options{
|
||||
Repo: repo, Store: store, Registry: registry.New(), CacheDir: t.TempDir(), Activators: []Activator{act},
|
||||
})
|
||||
if err := single.Reconcile(ctx); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
instances, err := single.Driver(manifest.RuntimeHost).Status(ctx, "acme.kit")
|
||||
if err != nil || len(instances) != 1 || instances[0].MemoryBytes != 12<<20 {
|
||||
t.Fatalf("instances = %+v, %v", instances, err)
|
||||
}
|
||||
act.usage = Usage{Idle: true}
|
||||
if s, _ := single.Status("acme.kit"); !s.Idle || s.MemoryBytes != 0 {
|
||||
t.Fatalf("status = %+v", s)
|
||||
}
|
||||
|
||||
mr := miniredis.RunT(t)
|
||||
rdb := redis.NewClient(&redis.Options{Addr: mr.Addr()})
|
||||
newNode := func(role string, u Usage) *Reconciler {
|
||||
return New(Options{
|
||||
Repo: repo, Store: store, Registry: registry.New(), CacheDir: t.TempDir(), Redis: rdb,
|
||||
Interval: time.Hour, Role: role, Activators: []Activator{&usageActivator{usage: u}},
|
||||
})
|
||||
}
|
||||
a, b := newNode("a", Usage{MemoryBytes: 10 << 20}), newNode("b", Usage{Idle: true})
|
||||
for _, r := range []*Reconciler{a, b} {
|
||||
if err := r.Reconcile(ctx); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
instances, err = a.Driver(manifest.RuntimeHost).Status(ctx, "acme.kit")
|
||||
if err != nil || len(instances) != 2 {
|
||||
t.Fatalf("instances = %+v, %v", instances, err)
|
||||
}
|
||||
got := map[string]driver.InstanceStatus{}
|
||||
for _, in := range instances {
|
||||
got[in.Node[:1]] = in
|
||||
}
|
||||
if got["a"].MemoryBytes != 10<<20 || got["a"].Idle || !got["b"].Idle {
|
||||
t.Fatalf("usage by node = %+v", got)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -150,6 +150,23 @@ type EgressReporter interface {
|
||||
Egress(pluginID string) driver.EgressMode
|
||||
}
|
||||
|
||||
// UsageReporter is an Activator that runs plugin code and can say what a
|
||||
// plugin's processes use on this node, for the node's status report.
|
||||
type UsageReporter interface {
|
||||
Activator
|
||||
// Usage is false for a plugin the activator does not run here.
|
||||
Usage(pluginID string) (Usage, bool)
|
||||
}
|
||||
|
||||
// Usage is what a plugin's processes use on one node.
|
||||
type Usage struct {
|
||||
// MemoryBytes is the resident memory of the plugin's processes; 0 when
|
||||
// not measured (the plugin is starting, or the system cannot tell).
|
||||
MemoryBytes int64 `json:"memoryBytes,omitempty"`
|
||||
// Idle: stopped for going without calls; the next call starts it.
|
||||
Idle bool `json:"idle,omitempty"`
|
||||
}
|
||||
|
||||
// PendingError is returned by an activator that started the plugin but
|
||||
// finishes in the background, such as a kubernetes rollout: the plugin is
|
||||
// loaded and shows as degraded with the reason until the activator reports
|
||||
@@ -178,6 +195,8 @@ type Status struct {
|
||||
UpgradeVersion string `json:"upgradeVersion,omitempty"`
|
||||
UpgradeState string `json:"upgradeState,omitempty"`
|
||||
UpgradeError string `json:"upgradeError,omitempty"`
|
||||
// Usage is what the plugin's processes use here, as of the report.
|
||||
Usage
|
||||
}
|
||||
|
||||
// Upgrade states reported in Status.
|
||||
@@ -825,11 +844,26 @@ func (r *Reconciler) ReportRuntime(pluginID string, healthy bool, err error) {
|
||||
// Status reports how one plugin fares on this node.
|
||||
func (r *Reconciler) Status(pluginID string) (Status, bool) {
|
||||
r.statusMu.RLock()
|
||||
defer r.statusMu.RUnlock()
|
||||
s, ok := r.status[pluginID]
|
||||
r.statusMu.RUnlock()
|
||||
if ok {
|
||||
s.Usage = r.usage(pluginID)
|
||||
}
|
||||
return s, ok
|
||||
}
|
||||
|
||||
// usage asks the activators what a plugin's processes use on this node.
|
||||
func (r *Reconciler) usage(pluginID string) Usage {
|
||||
for _, a := range r.activators {
|
||||
if u, ok := a.(UsageReporter); ok {
|
||||
if got, ok := u.Usage(pluginID); ok {
|
||||
return got
|
||||
}
|
||||
}
|
||||
}
|
||||
return Usage{}
|
||||
}
|
||||
|
||||
// Loaded returns the plugins loaded on this node, sorted by ID.
|
||||
func (r *Reconciler) Loaded() []*Loaded {
|
||||
r.mu.Lock()
|
||||
|
||||
@@ -160,6 +160,9 @@ weknora-plugin verify -pubkey ed25519:... acme-search-1.0.0.wkp
|
||||
- 默认不回收;桌面端默认 `10m`。
|
||||
- 停掉的插件仍在「插件管理」中显示为运行中,路由、独立宿主的通告都不受影响;升级时照常先启动新版本校验,再按空闲时间回收。
|
||||
- 插件在两次调用之间要做事(如后台轮询、长连接)时,在 `plugin.yaml` 中声明 `runtime.keepAlive: true`,始终常驻;声明了 `singleton` 的插件同样常驻。
|
||||
- 内存占用:「插件管理」列表的「内存」列显示插件各实例的常驻内存(RSS)合计,插件详情的「节点状态」逐个节点显示;被空闲回收的显示「空闲」。统计的是插件进程所在的整个进程组,包括插件派生的子进程和 Linux 网络沙箱的转发进程。
|
||||
- Linux 读 `/proc`,macOS 等系统每次采样运行一次 `ps`;Windows 不统计。采样结果缓存 10 秒,各节点随状态每 30 秒上报一次。
|
||||
- 远程插件与 Kubernetes 插件不统计,显示为「—」。
|
||||
- 在 Linux 上,插件进程按 `runtime.resources` 限制资源:
|
||||
- 内存总是受限。插件没有声明时,按 `WEKNORA_PLUGIN_MEMORY_DEFAULT`(如 `1Gi`)限制,未设置则不限。
|
||||
- CPU 需要把一个可写的 cgroup v2 目录委托给 WeKnora,并通过 `WEKNORA_PLUGIN_CGROUP` 指定。设置后,每个插件进程进入各自的子 cgroup。
|
||||
|
||||
Reference in New Issue
Block a user