diff --git a/go.mod b/go.mod index 72b1fb7..56ab4bb 100644 --- a/go.mod +++ b/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 diff --git a/go.sum b/go.sum new file mode 100644 index 0000000..dd5797d --- /dev/null +++ b/go.sum @@ -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= diff --git a/internal/checkpoint/file.go b/internal/checkpoint/file.go new file mode 100644 index 0000000..78ea6cf --- /dev/null +++ b/internal/checkpoint/file.go @@ -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 } diff --git a/internal/checkpoint/memory.go b/internal/checkpoint/memory.go new file mode 100644 index 0000000..96bdb42 --- /dev/null +++ b/internal/checkpoint/memory.go @@ -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 } diff --git a/internal/checkpoint/redis.go b/internal/checkpoint/redis.go new file mode 100644 index 0000000..c5187d3 --- /dev/null +++ b/internal/checkpoint/redis.go @@ -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() } diff --git a/internal/collect/collector.go b/internal/collect/collector.go new file mode 100644 index 0000000..c98fa55 --- /dev/null +++ b/internal/collect/collector.go @@ -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 +} diff --git a/internal/collect/collector_test.go b/internal/collect/collector_test.go new file mode 100644 index 0000000..fb56c55 --- /dev/null +++ b/internal/collect/collector_test.go @@ -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) + } + } +} diff --git a/internal/collect/logtail.go b/internal/collect/logtail.go new file mode 100644 index 0000000..6771fff --- /dev/null +++ b/internal/collect/logtail.go @@ -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" +} diff --git a/internal/collect/system.go b/internal/collect/system.go new file mode 100644 index 0000000..c3c1b01 --- /dev/null +++ b/internal/collect/system.go @@ -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 +} diff --git a/internal/config/config.go b/internal/config/config.go index 7e9bd49..760ee2e 100644 --- a/internal/config/config.go +++ b/internal/config/config.go @@ -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", "")),