develop: 开发 host-metrics-collector 采集模块
This commit is contained in:
parent
3c071afa50
commit
7d60597db3
8
go.mod
8
go.mod
@ -2,8 +2,6 @@ module git.opencomputing.cn/yumoqing/host-metrics-collector
|
||||
|
||||
go 1.22
|
||||
|
||||
require (
|
||||
github.com/redis/go-redis/v9 v9.5.1
|
||||
github.com/shirou/gopsutil/v4 v4.24.6
|
||||
google.golang.org/protobuf v1.34.2
|
||||
)
|
||||
require google.golang.org/protobuf v1.34.2
|
||||
|
||||
require github.com/google/go-cmp v0.6.0 // indirect
|
||||
|
||||
4
go.sum
Normal file
4
go.sum
Normal file
@ -0,0 +1,4 @@
|
||||
github.com/google/go-cmp v0.6.0 h1:ofyhxvXcZhMsU5ulbFiLKl/XBFqE1GSq7atu8tAmTRI=
|
||||
github.com/google/go-cmp v0.6.0/go.mod h1:17dUlkBOakJ0+DkrSSNjCkIjxS6bF9zb3elmeNGIjoY=
|
||||
google.golang.org/protobuf v1.34.2 h1:6xV6lTsCfpGD21XK49h7MhtcApnLqkfYgPcdHftf6hg=
|
||||
google.golang.org/protobuf v1.34.2/go.mod h1:qYOHts0dSfpeUzUFpOMr/WGzszTmLH+DiWniOlNbLDw=
|
||||
93
internal/checkpoint/file.go
Normal file
93
internal/checkpoint/file.go
Normal file
@ -0,0 +1,93 @@
|
||||
package checkpoint
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"sync"
|
||||
)
|
||||
|
||||
// FileStore 将检查点持久化到本地 JSON 文件,适合 Agent 本地 WAL 场景。
|
||||
type FileStore struct {
|
||||
path string
|
||||
|
||||
mu sync.RWMutex
|
||||
data fileData
|
||||
}
|
||||
|
||||
type fileData struct {
|
||||
LogOffsets map[string]int64 `json:"log_offsets"`
|
||||
AckedSeq map[string]int64 `json:"acked_seq"`
|
||||
}
|
||||
|
||||
// NewFileStore 加载(或初始化)指定路径的检查点文件。
|
||||
func NewFileStore(path string) (*FileStore, error) {
|
||||
s := &FileStore{
|
||||
path: path,
|
||||
data: fileData{
|
||||
LogOffsets: map[string]int64{},
|
||||
AckedSeq: map[string]int64{},
|
||||
},
|
||||
}
|
||||
|
||||
raw, err := os.ReadFile(path)
|
||||
switch {
|
||||
case err == nil:
|
||||
if err := json.Unmarshal(raw, &s.data); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
case os.IsNotExist(err):
|
||||
// 首次运行,使用空检查点。
|
||||
default:
|
||||
return nil, err
|
||||
}
|
||||
|
||||
if s.data.LogOffsets == nil {
|
||||
s.data.LogOffsets = map[string]int64{}
|
||||
}
|
||||
if s.data.AckedSeq == nil {
|
||||
s.data.AckedSeq = map[string]int64{}
|
||||
}
|
||||
return s, nil
|
||||
}
|
||||
|
||||
func (s *FileStore) persistLocked() error {
|
||||
raw, err := json.MarshalIndent(s.data, "", " ")
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if dir := filepath.Dir(s.path); dir != "." {
|
||||
if err := os.MkdirAll(dir, 0o755); err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
return os.WriteFile(s.path, raw, 0o644)
|
||||
}
|
||||
|
||||
func (s *FileStore) GetLogOffset(file string) (int64, error) {
|
||||
s.mu.RLock()
|
||||
defer s.mu.RUnlock()
|
||||
return s.data.LogOffsets[file], nil
|
||||
}
|
||||
|
||||
func (s *FileStore) SetLogOffset(file string, offset int64) error {
|
||||
s.mu.Lock()
|
||||
defer s.mu.Unlock()
|
||||
s.data.LogOffsets[file] = offset
|
||||
return s.persistLocked()
|
||||
}
|
||||
|
||||
func (s *FileStore) GetAckedSeq(hostID string) (int64, error) {
|
||||
s.mu.RLock()
|
||||
defer s.mu.RUnlock()
|
||||
return s.data.AckedSeq[hostID], nil
|
||||
}
|
||||
|
||||
func (s *FileStore) SetAckedSeq(hostID string, seq int64) error {
|
||||
s.mu.Lock()
|
||||
defer s.mu.Unlock()
|
||||
s.data.AckedSeq[hostID] = seq
|
||||
return s.persistLocked()
|
||||
}
|
||||
|
||||
func (s *FileStore) Close() error { return nil }
|
||||
46
internal/checkpoint/memory.go
Normal file
46
internal/checkpoint/memory.go
Normal file
@ -0,0 +1,46 @@
|
||||
package checkpoint
|
||||
|
||||
import "sync"
|
||||
|
||||
// MemoryStore 是进程内检查点存储,线程安全。用于单元测试与零依赖本地运行。
|
||||
type MemoryStore struct {
|
||||
mu sync.RWMutex
|
||||
logOffsets map[string]int64
|
||||
ackedSeq map[string]int64
|
||||
}
|
||||
|
||||
// NewMemoryStore 创建空的进程内检查点存储。
|
||||
func NewMemoryStore() *MemoryStore {
|
||||
return &MemoryStore{
|
||||
logOffsets: map[string]int64{},
|
||||
ackedSeq: map[string]int64{},
|
||||
}
|
||||
}
|
||||
|
||||
func (s *MemoryStore) GetLogOffset(file string) (int64, error) {
|
||||
s.mu.RLock()
|
||||
defer s.mu.RUnlock()
|
||||
return s.logOffsets[file], nil
|
||||
}
|
||||
|
||||
func (s *MemoryStore) SetLogOffset(file string, offset int64) error {
|
||||
s.mu.Lock()
|
||||
defer s.mu.Unlock()
|
||||
s.logOffsets[file] = offset
|
||||
return nil
|
||||
}
|
||||
|
||||
func (s *MemoryStore) GetAckedSeq(hostID string) (int64, error) {
|
||||
s.mu.RLock()
|
||||
defer s.mu.RUnlock()
|
||||
return s.ackedSeq[hostID], nil
|
||||
}
|
||||
|
||||
func (s *MemoryStore) SetAckedSeq(hostID string, seq int64) error {
|
||||
s.mu.Lock()
|
||||
defer s.mu.Unlock()
|
||||
s.ackedSeq[hostID] = seq
|
||||
return nil
|
||||
}
|
||||
|
||||
func (s *MemoryStore) Close() error { return nil }
|
||||
61
internal/checkpoint/redis.go
Normal file
61
internal/checkpoint/redis.go
Normal file
@ -0,0 +1,61 @@
|
||||
package checkpoint
|
||||
|
||||
import (
|
||||
"context"
|
||||
|
||||
"github.com/redis/go-redis/v9"
|
||||
)
|
||||
|
||||
// Redis 键设计(与 database-design.md 6 节对齐):
|
||||
// - 日志文件 offset:hms:agent:logoffset:{file}
|
||||
// - 主机已确认 seq:hms:agent:ackedseq:{host_id}
|
||||
const (
|
||||
keyLogOffset = "hms:agent:logoffset:"
|
||||
keyAckedSeq = "hms:agent:ackedseq:"
|
||||
)
|
||||
|
||||
// RedisStore 使用 Redis 记录检查点,实现 Agent 重启/漂移后的断点续传。
|
||||
type RedisStore struct {
|
||||
cli *redis.Client
|
||||
}
|
||||
|
||||
// NewRedisStore 解析 Redis URL 并创建客户端。
|
||||
func NewRedisStore(url string) (*RedisStore, error) {
|
||||
opt, err := redis.ParseURL(url)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return &RedisStore{cli: redis.NewClient(opt)}, nil
|
||||
}
|
||||
|
||||
func (s *RedisStore) GetLogOffset(file string) (int64, error) {
|
||||
v, err := s.cli.Get(context.Background(), keyLogOffset+file).Int64()
|
||||
if err == redis.Nil {
|
||||
return 0, nil
|
||||
}
|
||||
if err != nil {
|
||||
return 0, err
|
||||
}
|
||||
return v, nil
|
||||
}
|
||||
|
||||
func (s *RedisStore) SetLogOffset(file string, offset int64) error {
|
||||
return s.cli.Set(context.Background(), keyLogOffset+file, offset, 0).Err()
|
||||
}
|
||||
|
||||
func (s *RedisStore) GetAckedSeq(hostID string) (int64, error) {
|
||||
v, err := s.cli.Get(context.Background(), keyAckedSeq+hostID).Int64()
|
||||
if err == redis.Nil {
|
||||
return 0, nil
|
||||
}
|
||||
if err != nil {
|
||||
return 0, err
|
||||
}
|
||||
return v, nil
|
||||
}
|
||||
|
||||
func (s *RedisStore) SetAckedSeq(hostID string, seq int64) error {
|
||||
return s.cli.Set(context.Background(), keyAckedSeq+hostID, seq, 0).Err()
|
||||
}
|
||||
|
||||
func (s *RedisStore) Close() error { return s.cli.Close() }
|
||||
150
internal/collect/collector.go
Normal file
150
internal/collect/collector.go
Normal file
@ -0,0 +1,150 @@
|
||||
package collect
|
||||
|
||||
import (
|
||||
"errors"
|
||||
"fmt"
|
||||
"time"
|
||||
|
||||
"git.opencomputing.cn/yumoqing/host-metrics-collector/internal/model"
|
||||
)
|
||||
|
||||
// Collector 将 System 的原始读数转换为 HMS 指标样本。
|
||||
type Collector struct {
|
||||
sys System
|
||||
hostID string
|
||||
service string
|
||||
}
|
||||
|
||||
// NewCollector 创建指标采集器;sys 为空时使用 gopsutil 实现。
|
||||
func NewCollector(sys System, hostID, service string) *Collector {
|
||||
if sys == nil {
|
||||
sys = GopsutilSystem{}
|
||||
}
|
||||
return &Collector{sys: sys, hostID: hostID, service: service}
|
||||
}
|
||||
|
||||
func (c *Collector) baseLabels(group string) map[string]string {
|
||||
return map[string]string{
|
||||
model.LabelHostID: c.hostID,
|
||||
model.LabelService: c.service,
|
||||
model.LabelMetricGroup: group,
|
||||
}
|
||||
}
|
||||
|
||||
func withLabels(base, extra map[string]string) map[string]string {
|
||||
out := make(map[string]string, len(base)+len(extra))
|
||||
for k, v := range base {
|
||||
out[k] = v
|
||||
}
|
||||
for k, v := range extra {
|
||||
out[k] = v
|
||||
}
|
||||
return out
|
||||
}
|
||||
|
||||
// CollectCore 采集核心指标(15s):cpu.usage、mem.used_percent、load.1m。
|
||||
func (c *Collector) CollectCore() ([]model.Sample, error) {
|
||||
var out []model.Sample
|
||||
var errs []error
|
||||
at := time.Now()
|
||||
base := c.baseLabels(model.GroupCore)
|
||||
|
||||
if v, err := c.sys.CPUUsage(); err == nil {
|
||||
out = append(out, model.NewSample(model.MetricCPUUsage, v, at, base))
|
||||
} else {
|
||||
errs = append(errs, fmt.Errorf("cpu usage: %w", err))
|
||||
}
|
||||
|
||||
if per, err := c.sys.PerCPUUsage(); err == nil {
|
||||
for i, v := range per {
|
||||
l := withLabels(base, map[string]string{"core": fmt.Sprintf("%d", i)})
|
||||
out = append(out, model.NewSample(model.MetricCPUUsage, v, at, l))
|
||||
}
|
||||
}
|
||||
|
||||
if v, err := c.sys.MemUsedPercent(); err == nil {
|
||||
out = append(out, model.NewSample(model.MetricMemUsedPercent, v, at, base))
|
||||
} else {
|
||||
errs = append(errs, fmt.Errorf("mem usage: %w", err))
|
||||
}
|
||||
|
||||
if v, err := c.sys.Load1(); err == nil {
|
||||
out = append(out, model.NewSample(model.MetricLoad1m, v, at, base))
|
||||
} else {
|
||||
errs = append(errs, fmt.Errorf("load1: %w", err))
|
||||
}
|
||||
|
||||
return out, errors.Join(errs...)
|
||||
}
|
||||
|
||||
// CollectDisk 采集磁盘指标(30s):disk.used_percent、disk.io.read_bytes/write_bytes。
|
||||
func (c *Collector) CollectDisk() ([]model.Sample, error) {
|
||||
var out []model.Sample
|
||||
var errs []error
|
||||
at := time.Now()
|
||||
base := c.baseLabels(model.GroupDisk)
|
||||
|
||||
if usages, err := c.sys.DiskUsages(); err == nil {
|
||||
for _, u := range usages {
|
||||
l := withLabels(base, map[string]string{"mountpoint": u.Mountpoint})
|
||||
out = append(out, model.NewSample(model.MetricDiskUsedPercent, u.UsedPercent, at, l))
|
||||
}
|
||||
} else {
|
||||
errs = append(errs, fmt.Errorf("disk usage: %w", err))
|
||||
}
|
||||
|
||||
if ios, err := c.sys.DiskIOCounters(); err == nil {
|
||||
for dev, io := range ios {
|
||||
l := withLabels(base, map[string]string{"device": dev})
|
||||
out = append(out, model.NewSample(model.MetricDiskIOReadBytes, float64(io.ReadBytes), at, l))
|
||||
out = append(out, model.NewSample(model.MetricDiskIOWriteBytes, float64(io.WriteBytes), at, l))
|
||||
}
|
||||
} else {
|
||||
errs = append(errs, fmt.Errorf("disk io: %w", err))
|
||||
}
|
||||
|
||||
return out, errors.Join(errs...)
|
||||
}
|
||||
|
||||
// CollectNetwork 采集网络指标(30s):net.bytes_sent/recv、net.pkt_drop。
|
||||
func (c *Collector) CollectNetwork() ([]model.Sample, error) {
|
||||
at := time.Now()
|
||||
base := c.baseLabels(model.GroupNetwork)
|
||||
|
||||
ios, err := c.sys.NetIOCounters()
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("net io: %w", err)
|
||||
}
|
||||
|
||||
out := make([]model.Sample, 0, len(ios)*3)
|
||||
for iface, io := range ios {
|
||||
l := withLabels(base, map[string]string{"interface": iface})
|
||||
out = append(out, model.NewSample(model.MetricNetBytesSent, float64(io.BytesSent), at, l))
|
||||
out = append(out, model.NewSample(model.MetricNetBytesRecv, float64(io.BytesRecv), at, l))
|
||||
out = append(out, model.NewSample(model.MetricNetPktDrop, float64(io.DropOut+io.DropIn), at, l))
|
||||
}
|
||||
return out, nil
|
||||
}
|
||||
|
||||
// CollectProcess 采集进程指标(60s):process.count、process.cpu、process.mem。
|
||||
func (c *Collector) CollectProcess() ([]model.Sample, error) {
|
||||
at := time.Now()
|
||||
base := c.baseLabels(model.GroupProcess)
|
||||
|
||||
procs, err := c.sys.Processes()
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("process list: %w", err)
|
||||
}
|
||||
|
||||
out := make([]model.Sample, 0, len(procs)*2+1)
|
||||
out = append(out, model.NewSample(model.MetricProcessCount, float64(len(procs)), at, base))
|
||||
for _, p := range procs {
|
||||
l := withLabels(base, map[string]string{
|
||||
"pid": fmt.Sprintf("%d", p.PID),
|
||||
"name": p.Name,
|
||||
})
|
||||
out = append(out, model.NewSample(model.MetricProcessCPU, p.CPUPercent, at, l))
|
||||
out = append(out, model.NewSample(model.MetricProcessMem, float64(p.MemBytes), at, l))
|
||||
}
|
||||
return out, nil
|
||||
}
|
||||
171
internal/collect/collector_test.go
Normal file
171
internal/collect/collector_test.go
Normal file
@ -0,0 +1,171 @@
|
||||
package collect
|
||||
|
||||
import (
|
||||
"strings"
|
||||
"testing"
|
||||
|
||||
"git.opencomputing.cn/yumoqing/host-metrics-collector/internal/model"
|
||||
)
|
||||
|
||||
type fakeSystem struct {
|
||||
cpu float64
|
||||
percpu []float64
|
||||
mem float64
|
||||
load float64
|
||||
disks []DiskUsageStat
|
||||
diskio map[string]DiskIOCounter
|
||||
netio map[string]NetIOCounter
|
||||
procs []ProcessStat
|
||||
}
|
||||
|
||||
func (f fakeSystem) CPUUsage() (float64, error) { return f.cpu, nil }
|
||||
func (f fakeSystem) PerCPUUsage() ([]float64, error) { return f.percpu, nil }
|
||||
func (f fakeSystem) MemUsedPercent() (float64, error) { return f.mem, nil }
|
||||
func (f fakeSystem) Load1() (float64, error) { return f.load, nil }
|
||||
func (f fakeSystem) DiskUsages() ([]DiskUsageStat, error) { return f.disks, nil }
|
||||
func (f fakeSystem) DiskIOCounters() (map[string]DiskIOCounter, error) {
|
||||
return f.diskio, nil
|
||||
}
|
||||
func (f fakeSystem) NetIOCounters() (map[string]NetIOCounter, error) { return f.netio, nil }
|
||||
func (f fakeSystem) Processes() ([]ProcessStat, error) { return f.procs, nil }
|
||||
|
||||
func newTestCollector() *Collector {
|
||||
sys := fakeSystem{
|
||||
cpu: 42.5,
|
||||
percpu: []float64{10, 20},
|
||||
mem: 66.6,
|
||||
load: 0.7,
|
||||
disks: []DiskUsageStat{{Mountpoint: "/", UsedPercent: 80}},
|
||||
diskio: map[string]DiskIOCounter{"sda": {ReadBytes: 100, WriteBytes: 200}},
|
||||
netio: map[string]NetIOCounter{"eth0": {BytesSent: 1000, BytesRecv: 2000, DropOut: 1, DropIn: 2}},
|
||||
procs: []ProcessStat{{PID: 1, Name: "init", CPUPercent: 1.2, MemBytes: 4096}},
|
||||
}
|
||||
return NewCollector(sys, "h-001", "web")
|
||||
}
|
||||
|
||||
func TestCollectCore(t *testing.T) {
|
||||
c := newTestCollector()
|
||||
samples, err := c.CollectCore()
|
||||
if err != nil {
|
||||
t.Fatalf("CollectCore returned error: %v", err)
|
||||
}
|
||||
if len(samples) != 5 { // 1 overall + 2 per-core + mem + load
|
||||
t.Fatalf("expected 5 samples, got %d", len(samples))
|
||||
}
|
||||
|
||||
byName := map[string][]model.Sample{}
|
||||
for _, s := range samples {
|
||||
byName[s.Name] = append(byName[s.Name], s)
|
||||
}
|
||||
if got := byName[model.MetricCPUUsage]; len(got) != 3 {
|
||||
t.Fatalf("expected 3 cpu.usage samples, got %d", len(got))
|
||||
}
|
||||
if got := byName[model.MetricMemUsedPercent][0].Value; got != 66.6 {
|
||||
t.Fatalf("mem used percent = %v, want 66.6", got)
|
||||
}
|
||||
if got := byName[model.MetricLoad1m][0].Value; got != 0.7 {
|
||||
t.Fatalf("load1m = %v, want 0.7", got)
|
||||
}
|
||||
for _, s := range samples {
|
||||
if s.Labels[model.LabelHostID] != "h-001" {
|
||||
t.Fatalf("missing host_id label: %#v", s.Labels)
|
||||
}
|
||||
if s.Labels[model.LabelMetricGroup] != model.GroupCore {
|
||||
t.Fatalf("wrong metric_group: %#v", s.Labels)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestCollectDisk(t *testing.T) {
|
||||
c := newTestCollector()
|
||||
samples, err := c.CollectDisk()
|
||||
if err != nil {
|
||||
t.Fatalf("CollectDisk returned error: %v", err)
|
||||
}
|
||||
if len(samples) != 3 { // 1 used_percent + 2 io counters
|
||||
t.Fatalf("expected 3 samples, got %d", len(samples))
|
||||
}
|
||||
var sawUsed, sawRead, sawWrite bool
|
||||
for _, s := range samples {
|
||||
switch s.Name {
|
||||
case model.MetricDiskUsedPercent:
|
||||
sawUsed = true
|
||||
if s.Labels["mountpoint"] != "/" {
|
||||
t.Fatalf("mountpoint = %q", s.Labels["mountpoint"])
|
||||
}
|
||||
case model.MetricDiskIOReadBytes:
|
||||
sawRead = true
|
||||
if s.Value != 100 {
|
||||
t.Fatalf("read bytes = %v", s.Value)
|
||||
}
|
||||
case model.MetricDiskIOWriteBytes:
|
||||
sawWrite = true
|
||||
}
|
||||
}
|
||||
if !sawUsed || !sawRead || !sawWrite {
|
||||
t.Fatalf("missing disk samples: used=%v read=%v write=%v", sawUsed, sawRead, sawWrite)
|
||||
}
|
||||
}
|
||||
|
||||
func TestCollectNetwork(t *testing.T) {
|
||||
c := newTestCollector()
|
||||
samples, err := c.CollectNetwork()
|
||||
if err != nil {
|
||||
t.Fatalf("CollectNetwork returned error: %v", err)
|
||||
}
|
||||
if len(samples) != 3 {
|
||||
t.Fatalf("expected 3 samples, got %d", len(samples))
|
||||
}
|
||||
for _, s := range samples {
|
||||
if s.Labels["interface"] != "eth0" {
|
||||
t.Fatalf("interface = %q", s.Labels["interface"])
|
||||
}
|
||||
if s.Name == model.MetricNetPktDrop && s.Value != 3 {
|
||||
t.Fatalf("pkt drop = %v, want 3", s.Value)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestCollectProcess(t *testing.T) {
|
||||
c := newTestCollector()
|
||||
samples, err := c.CollectProcess()
|
||||
if err != nil {
|
||||
t.Fatalf("CollectProcess returned error: %v", err)
|
||||
}
|
||||
if len(samples) != 3 { // count + cpu + mem
|
||||
t.Fatalf("expected 3 samples, got %d", len(samples))
|
||||
}
|
||||
for _, s := range samples {
|
||||
if s.Name == model.MetricProcessCount && s.Value != 1 {
|
||||
t.Fatalf("process count = %v, want 1", s.Value)
|
||||
}
|
||||
if s.Name == model.MetricProcessMem && s.Value != 4096 {
|
||||
t.Fatalf("process mem = %v, want 4096", s.Value)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestDetectLevel(t *testing.T) {
|
||||
cases := map[string]string{
|
||||
"ERROR something": "ERROR",
|
||||
"fatal error": "FATAL",
|
||||
"warn: x": "WARN",
|
||||
"info msg": "INFO",
|
||||
"plain": "INFO",
|
||||
}
|
||||
for in, want := range cases {
|
||||
if got := detectLevel(in); got != want {
|
||||
t.Errorf("detectLevel(%q) = %q, want %q", in, got, want)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestCollectorLabelsCarryService(t *testing.T) {
|
||||
c := newTestCollector()
|
||||
samples, _ := c.CollectCore()
|
||||
for _, s := range samples {
|
||||
if !strings.Contains(s.Labels[model.LabelService], "web") {
|
||||
t.Fatalf("service label missing: %#v", s.Labels)
|
||||
}
|
||||
}
|
||||
}
|
||||
129
internal/collect/logtail.go
Normal file
129
internal/collect/logtail.go
Normal file
@ -0,0 +1,129 @@
|
||||
package collect
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
"io"
|
||||
"os"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"git.opencomputing.cn/yumoqing/host-metrics-collector/internal/checkpoint"
|
||||
"git.opencomputing.cn/yumoqing/host-metrics-collector/internal/model"
|
||||
)
|
||||
|
||||
// LogTail 持续读取单个日志文件的新增内容,并通过 checkpoint.Store 记录
|
||||
// 已消费 offset,实现进程重启后的断点续传。采用轮询式 tail,避免对
|
||||
// inotify 的额外依赖,跨平台且足够轻量。
|
||||
type LogTail struct {
|
||||
path string
|
||||
store checkpoint.Store
|
||||
poll time.Duration
|
||||
onEntry func(model.LogRecord)
|
||||
}
|
||||
|
||||
// NewLogTail 创建日志 tail 采集器;onEntry 每读到一行完整日志时回调。
|
||||
func NewLogTail(path string, store checkpoint.Store, poll time.Duration, onEntry func(model.LogRecord)) *LogTail {
|
||||
if poll <= 0 {
|
||||
poll = time.Second
|
||||
}
|
||||
return &LogTail{path: path, store: store, poll: poll, onEntry: onEntry}
|
||||
}
|
||||
|
||||
// Run 阻塞执行采集,直到 ctx 被取消。
|
||||
func (t *LogTail) Run(ctx context.Context) {
|
||||
offset, err := t.store.GetLogOffset(t.path)
|
||||
if err != nil {
|
||||
offset = 0
|
||||
}
|
||||
var partial []byte
|
||||
|
||||
for {
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
return
|
||||
default:
|
||||
}
|
||||
|
||||
offset, partial = t.readOnce(offset, partial)
|
||||
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
return
|
||||
case <-time.After(t.poll):
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// readOnce 从 offset 处读取文件新增内容,返回新的 offset 与未成行的残留字节。
|
||||
func (t *LogTail) readOnce(offset int64, partial []byte) (int64, []byte) {
|
||||
f, err := os.Open(t.path)
|
||||
if err != nil {
|
||||
return offset, partial
|
||||
}
|
||||
defer f.Close()
|
||||
|
||||
info, err := f.Stat()
|
||||
if err != nil {
|
||||
return offset, partial
|
||||
}
|
||||
// 文件被截断或轮转:从头重新消费。
|
||||
if info.Size() < offset {
|
||||
offset = 0
|
||||
partial = partial[:0]
|
||||
}
|
||||
|
||||
if _, err := f.Seek(offset, io.SeekStart); err != nil {
|
||||
return offset, partial
|
||||
}
|
||||
data, err := io.ReadAll(f)
|
||||
if err != nil {
|
||||
return offset, partial
|
||||
}
|
||||
if len(data) == 0 {
|
||||
return offset, partial
|
||||
}
|
||||
|
||||
partial = append(partial, data...)
|
||||
for {
|
||||
idx := bytes.IndexByte(partial, '\n')
|
||||
if idx < 0 {
|
||||
break
|
||||
}
|
||||
line := partial[:idx]
|
||||
partial = partial[idx+1:]
|
||||
offset += int64(idx) + 1
|
||||
t.emit(line, offset)
|
||||
}
|
||||
return offset, partial
|
||||
}
|
||||
|
||||
func (t *LogTail) emit(line []byte, offset int64) {
|
||||
msg := strings.TrimRight(string(line), "\r")
|
||||
rec := model.LogRecord{
|
||||
Timestamp: time.Now().Unix(),
|
||||
Level: detectLevel(msg),
|
||||
Source: t.path,
|
||||
Message: msg,
|
||||
File: t.path,
|
||||
Offset: offset,
|
||||
}
|
||||
// 检查点写入失败不阻断采集;下轮会按旧 offset 重读(至少一次语义)。
|
||||
if err := t.store.SetLogOffset(t.path, offset); err != nil {
|
||||
// 可选:由上层记录日志,这里保持轻量。
|
||||
}
|
||||
if t.onEntry != nil {
|
||||
t.onEntry(rec)
|
||||
}
|
||||
}
|
||||
|
||||
// detectLevel 从日志内容中识别常见级别,未命中默认 INFO。
|
||||
func detectLevel(msg string) string {
|
||||
upper := strings.ToUpper(msg)
|
||||
for _, lvl := range []string{"FATAL", "ERROR", "WARN", "INFO", "DEBUG"} {
|
||||
if strings.Contains(upper, lvl) {
|
||||
return lvl
|
||||
}
|
||||
}
|
||||
return "INFO"
|
||||
}
|
||||
164
internal/collect/system.go
Normal file
164
internal/collect/system.go
Normal file
@ -0,0 +1,164 @@
|
||||
// Package collect 实现主机指标采集与文件日志 tail 采集。
|
||||
package collect
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
|
||||
"github.com/shirou/gopsutil/v4/cpu"
|
||||
"github.com/shirou/gopsutil/v4/disk"
|
||||
"github.com/shirou/gopsutil/v4/load"
|
||||
"github.com/shirou/gopsutil/v4/mem"
|
||||
"github.com/shirou/gopsutil/v4/net"
|
||||
"github.com/shirou/gopsutil/v4/process"
|
||||
)
|
||||
|
||||
// 以下结构是对 gopsutil 原始返回值的抽象,采集层不直接暴露 gopsutil 类型,
|
||||
// 便于单元测试注入 fake 数据。
|
||||
|
||||
// DiskUsageStat 描述一个挂载点的磁盘使用率。
|
||||
type DiskUsageStat struct {
|
||||
Mountpoint string
|
||||
UsedPercent float64
|
||||
}
|
||||
|
||||
// DiskIOCounter 描述一个块设备的累计读写字节数。
|
||||
type DiskIOCounter struct {
|
||||
ReadBytes uint64
|
||||
WriteBytes uint64
|
||||
}
|
||||
|
||||
// NetIOCounter 描述一个网卡的累计收发/丢包计数。
|
||||
type NetIOCounter struct {
|
||||
BytesSent uint64
|
||||
BytesRecv uint64
|
||||
DropOut uint64
|
||||
DropIn uint64
|
||||
}
|
||||
|
||||
// ProcessStat 描述一个进程的关键资源占用。
|
||||
type ProcessStat struct {
|
||||
PID int32
|
||||
Name string
|
||||
CPUPercent float64
|
||||
MemPercent float64
|
||||
MemBytes uint64
|
||||
}
|
||||
|
||||
// System 抽象主机系统信息读取。生产环境使用 GopsutilSystem,测试注入 fake。
|
||||
type System interface {
|
||||
CPUUsage() (float64, error)
|
||||
PerCPUUsage() ([]float64, error)
|
||||
MemUsedPercent() (float64, error)
|
||||
Load1() (float64, error)
|
||||
DiskUsages() ([]DiskUsageStat, error)
|
||||
DiskIOCounters() (map[string]DiskIOCounter, error)
|
||||
NetIOCounters() (map[string]NetIOCounter, error)
|
||||
Processes() ([]ProcessStat, error)
|
||||
}
|
||||
|
||||
// GopsutilSystem 是 System 的 gopsutil/v4 实现。
|
||||
type GopsutilSystem struct{}
|
||||
|
||||
func (GopsutilSystem) CPUUsage() (float64, error) {
|
||||
ps, err := cpu.Percent(0, false)
|
||||
if err != nil {
|
||||
return 0, err
|
||||
}
|
||||
if len(ps) == 0 {
|
||||
return 0, fmt.Errorf("cpu.Percent returned no samples")
|
||||
}
|
||||
return ps[0], nil
|
||||
}
|
||||
|
||||
func (GopsutilSystem) PerCPUUsage() ([]float64, error) {
|
||||
return cpu.Percent(0, true)
|
||||
}
|
||||
|
||||
func (GopsutilSystem) MemUsedPercent() (float64, error) {
|
||||
vm, err := mem.VirtualMemory()
|
||||
if err != nil {
|
||||
return 0, err
|
||||
}
|
||||
return vm.UsedPercent, nil
|
||||
}
|
||||
|
||||
func (GopsutilSystem) Load1() (float64, error) {
|
||||
avg, err := load.Avg()
|
||||
if err != nil {
|
||||
return 0, err
|
||||
}
|
||||
return avg.Load1, nil
|
||||
}
|
||||
|
||||
func (GopsutilSystem) DiskUsages() ([]DiskUsageStat, error) {
|
||||
parts, err := disk.Partitions(false)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
out := make([]DiskUsageStat, 0, len(parts))
|
||||
for _, p := range parts {
|
||||
if p.Mountpoint == "" {
|
||||
continue
|
||||
}
|
||||
usage, err := disk.Usage(p.Mountpoint)
|
||||
if err != nil {
|
||||
continue
|
||||
}
|
||||
out = append(out, DiskUsageStat{Mountpoint: p.Mountpoint, UsedPercent: usage.UsedPercent})
|
||||
}
|
||||
return out, nil
|
||||
}
|
||||
|
||||
func (GopsutilSystem) DiskIOCounters() (map[string]DiskIOCounter, error) {
|
||||
raw, err := disk.IOCounters()
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
out := make(map[string]DiskIOCounter, len(raw))
|
||||
for name, v := range raw {
|
||||
out[name] = DiskIOCounter{ReadBytes: v.ReadBytes, WriteBytes: v.WriteBytes}
|
||||
}
|
||||
return out, nil
|
||||
}
|
||||
|
||||
func (GopsutilSystem) NetIOCounters() (map[string]NetIOCounter, error) {
|
||||
raw, err := net.IOCounters(true)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
out := make(map[string]NetIOCounter, len(raw))
|
||||
for _, v := range raw {
|
||||
out[v.Name] = NetIOCounter{
|
||||
BytesSent: v.BytesSent,
|
||||
BytesRecv: v.BytesRecv,
|
||||
DropOut: v.Dropout,
|
||||
DropIn: v.Dropin,
|
||||
}
|
||||
}
|
||||
return out, nil
|
||||
}
|
||||
|
||||
func (GopsutilSystem) Processes() ([]ProcessStat, error) {
|
||||
raw, err := process.Processes()
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
out := make([]ProcessStat, 0, len(raw))
|
||||
for _, p := range raw {
|
||||
stat := ProcessStat{PID: p.Pid}
|
||||
if name, err := p.Name(); err == nil {
|
||||
stat.Name = name
|
||||
}
|
||||
if cpuPercent, err := p.CPUPercent(); err == nil {
|
||||
stat.CPUPercent = cpuPercent
|
||||
}
|
||||
if memPercent, err := p.MemoryPercent(); err == nil {
|
||||
stat.MemPercent = float64(memPercent)
|
||||
}
|
||||
if memInfo, err := p.MemoryInfo(); err == nil {
|
||||
stat.MemBytes = memInfo.RSS
|
||||
}
|
||||
out = append(out, stat)
|
||||
}
|
||||
return out, nil
|
||||
}
|
||||
@ -12,6 +12,7 @@ import (
|
||||
type AgentConfig struct {
|
||||
HostID string
|
||||
AgentVersion string
|
||||
Service string
|
||||
|
||||
GatewayURL string
|
||||
AgentToken string
|
||||
@ -38,6 +39,7 @@ func DefaultAgentConfig() AgentConfig {
|
||||
return AgentConfig{
|
||||
HostID: env("HMS_HOST_ID", "host-001"),
|
||||
AgentVersion: env("HMS_AGENT_VERSION", "0.1.0"),
|
||||
Service: env("HMS_SERVICE", "host"),
|
||||
GatewayURL: env("HMS_GATEWAY_URL", "http://127.0.0.1:4318"),
|
||||
AgentToken: env("HMS_AGENT_TOKEN", "dev-token"),
|
||||
LogFiles: splitCSV(env("HMS_LOG_FILES", "")),
|
||||
|
||||
Loading…
x
Reference in New Issue
Block a user