From 4f3e335f2862f255077278abd67673360d2605fa Mon Sep 17 00:00:00 2001 From: Pipeline Agent Date: Sat, 15 Aug 2026 11:29:17 +0800 Subject: [PATCH] =?UTF-8?q?develop:=20=E9=87=8D=E5=81=9A=20host-metrics-co?= =?UTF-8?q?llector=20=E9=87=87=E9=9B=86=E6=A8=A1=E5=9D=97=EF=BC=88?= =?UTF-8?q?=E9=81=B5=E5=AE=88=E6=A8=A1=E5=9D=97=E5=BC=80=E5=8F=91=E8=A7=84?= =?UTF-8?q?=E8=8C=83=EF=BC=89?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- internal/agent/agent_test.go | 200 ++++++++++++++++++++++++++++++++ internal/agent/reporter_test.go | 167 ++++++++++++++++++++++++++ 2 files changed, 367 insertions(+) create mode 100644 internal/agent/agent_test.go create mode 100644 internal/agent/reporter_test.go diff --git a/internal/agent/agent_test.go b/internal/agent/agent_test.go new file mode 100644 index 0000000..3123105 --- /dev/null +++ b/internal/agent/agent_test.go @@ -0,0 +1,200 @@ +package agent + +import ( + "context" + "io" + "net/http" + "net/http/httptest" + "path/filepath" + "testing" + "time" + + "google.golang.org/protobuf/proto" + + "git.opencomputing.cn/yumoqing/host-metrics-collector/internal/checkpoint" + "git.opencomputing.cn/yumoqing/host-metrics-collector/internal/config" + "git.opencomputing.cn/yumoqing/host-metrics-collector/internal/model" + "git.opencomputing.cn/yumoqing/host-metrics-collector/internal/wal" + "git.opencomputing.cn/yumoqing/host-metrics-collector/pkg/collectorpb" +) + +func TestNextSeqMonotonic(t *testing.T) { + a := &Agent{} + seen := map[int64]bool{} + for i := 0; i < 100; i++ { + seq := a.nextMetricsSeq() + if seen[seq] { + t.Fatalf("duplicate metrics seq %d", seq) + } + seen[seq] = true + } + for i := 0; i < 100; i++ { + seq := a.nextLogsSeq() + if seq <= 0 { + t.Fatalf("logs seq = %d", seq) + } + } +} + +func TestInitSeqRestoresFromWAL(t *testing.T) { + dir := t.TempDir() + mw, err := wal.Open(filepath.Join(dir, "m.wal")) + if err != nil { + t.Fatalf("open metrics wal: %v", err) + } + defer mw.Close() + lw, err := wal.Open(filepath.Join(dir, "l.wal")) + if err != nil { + t.Fatalf("open logs wal: %v", err) + } + defer lw.Close() + + // 模拟进程重启前已写入但尚未确认的批次。 + if err := mw.Append(3, []byte("m3")); err != nil { + t.Fatalf("append m3: %v", err) + } + if err := mw.Append(5, []byte("m5")); err != nil { + t.Fatalf("append m5: %v", err) + } + if err := lw.Append(2, []byte("l2")); err != nil { + t.Fatalf("append l2: %v", err) + } + + store := checkpoint.NewMemoryStore() + // 已确认 seq 为 2,WAL 中还有 seq 3/5 未确认。 + if err := store.SetAckedSeq(metricsSeqKey+"h-001", 2); err != nil { + t.Fatalf("set acked: %v", err) + } + + a := &Agent{ + cfg: config.AgentConfig{HostID: "h-001"}, + store: store, + metricsWAL: mw, + logsWAL: lw, + } + if err := a.initSeq(); err != nil { + t.Fatalf("initSeq: %v", err) + } + if a.metricsSeq != 5 { + t.Fatalf("metricsSeq = %d, want 5", a.metricsSeq) + } + if a.logsSeq != 2 { + t.Fatalf("logsSeq = %d, want 2", a.logsSeq) + } +} + +// fakeGateway 返回一个接收 protobuf 上报并回 Ack 的 httptest 服务。 +func fakeGateway(t *testing.T) (*httptest.Server, *int64, *[]string) { + t.Helper() + var ackedSeq int64 + var metrics []string + ts := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + if r.Header.Get(headerAgentToken) != "test-token" { + http.Error(w, "unauthorized", http.StatusUnauthorized) + return + } + body, err := io.ReadAll(r.Body) + if err != nil { + http.Error(w, "read", http.StatusBadRequest) + return + } + var batch collectorpb.MetricBatch + if err := proto.Unmarshal(body, &batch); err != nil { + http.Error(w, "bad proto", http.StatusBadRequest) + return + } + for _, s := range batch.Samples { + metrics = append(metrics, s.Name) + } + ackedSeq = batch.Seq + w.Header().Set("Content-Type", contentTypeProtobuf) + data, _ := proto.Marshal(&collectorpb.Ack{Ok: true, Seq: batch.Seq}) + _, _ = w.Write(data) + })) + return ts, &ackedSeq, &metrics +} + +func TestFlushMetricsEndToEnd(t *testing.T) { + ts, ackedSeq, metrics := fakeGateway(t) + defer ts.Close() + + dir := t.TempDir() + cfg := config.DefaultAgentConfig() + cfg.GatewayURL = ts.URL + cfg.AgentToken = "test-token" + cfg.HostID = "h-001" + cfg.CheckpointDriver = "memory" + cfg.CheckpointFile = filepath.Join(dir, "cp") + + a, err := New(cfg) + if err != nil { + t.Fatalf("New: %v", err) + } + defer a.Close() + + at := time.Now() + samples := []model.Sample{ + model.NewSample(model.MetricCPUUsage, 42.5, at, map[string]string{model.LabelHostID: "h-001"}), + model.NewSample(model.MetricMemUsedPercent, 66.6, at, map[string]string{model.LabelHostID: "h-001"}), + } + if err := a.flushMetrics(context.Background(), samples); err != nil { + t.Fatalf("flushMetrics: %v", err) + } + + if got := *ackedSeq; got != 1 { + t.Fatalf("acked seq = %d, want 1", got) + } + if len(*metrics) != 2 { + t.Fatalf("metrics = %v", *metrics) + } + if got, _ := a.store.GetAckedSeq(metricsSeqKey + "h-001"); got != 1 { + t.Fatalf("local acked seq = %d, want 1", got) + } +} + +func TestFlushLogsEndToEnd(t *testing.T) { + var gotSeq int64 + var gotMsgs []string + ts := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + body, _ := io.ReadAll(r.Body) + var batch collectorpb.LogBatch + if err := proto.Unmarshal(body, &batch); err != nil { + t.Errorf("unmarshal: %v", err) + } + gotSeq = batch.Seq + for _, e := range batch.Entries { + gotMsgs = append(gotMsgs, e.Message) + } + w.Header().Set("Content-Type", contentTypeProtobuf) + data, _ := proto.Marshal(&collectorpb.Ack{Ok: true, Seq: batch.Seq}) + _, _ = w.Write(data) + })) + defer ts.Close() + + dir := t.TempDir() + cfg := config.DefaultAgentConfig() + cfg.GatewayURL = ts.URL + cfg.AgentToken = "test-token" + cfg.HostID = "h-001" + cfg.CheckpointDriver = "memory" + cfg.CheckpointFile = filepath.Join(dir, "cp") + + a, err := New(cfg) + if err != nil { + t.Fatalf("New: %v", err) + } + defer a.Close() + + records := []model.LogRecord{ + {Timestamp: time.Now().Unix(), Level: "ERROR", Message: "boom", File: "/var/log/app.log", Offset: 10}, + } + if err := a.flushLogs(context.Background(), records); err != nil { + t.Fatalf("flushLogs: %v", err) + } + if gotSeq != 1 || len(gotMsgs) != 1 || gotMsgs[0] != "boom" { + t.Fatalf("seq=%d msgs=%v", gotSeq, gotMsgs) + } + if got, _ := a.store.GetAckedSeq(logsSeqKey + "h-001"); got != 1 { + t.Fatalf("local acked seq = %d, want 1", got) + } +} diff --git a/internal/agent/reporter_test.go b/internal/agent/reporter_test.go new file mode 100644 index 0000000..d32ce2e --- /dev/null +++ b/internal/agent/reporter_test.go @@ -0,0 +1,167 @@ +package agent + +import ( + "context" + "io" + "net/http" + "net/http/httptest" + "strings" + "sync/atomic" + "testing" + "time" + + "google.golang.org/protobuf/proto" + + "git.opencomputing.cn/yumoqing/host-metrics-collector/pkg/collectorpb" +) + +func newTestReporter(url, token string) *Reporter { + r := NewReporter(url, token, "h-001", "0.1.0") + // 测试中缩短重试退避与尝试次数,避免拖慢用例。 + r.maxAttempts = 3 + r.baseBackoff = time.Millisecond + return r +} + +func ackProto(t *testing.T, ok bool, seq int64, msg string) []byte { + t.Helper() + data, err := proto.Marshal(&collectorpb.Ack{Ok: ok, Seq: seq, Error: msg}) + if err != nil { + t.Fatalf("marshal ack: %v", err) + } + return data +} + +func TestReporterSendMetricsSuccess(t *testing.T) { + var gotToken, gotHost, gotContentType string + ts := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + gotToken = r.Header.Get(headerAgentToken) + gotHost = r.Header.Get(headerHostID) + gotContentType = r.Header.Get("Content-Type") + w.Header().Set("Content-Type", contentTypeProtobuf) + _, _ = w.Write(ackProto(t, true, 42, "")) + })) + defer ts.Close() + + r := newTestReporter(ts.URL, "test-token") + ack, err := r.SendMetrics(context.Background(), &collectorpb.MetricBatch{HostId: "h-001", Seq: 42}) + if err != nil { + t.Fatalf("SendMetrics: %v", err) + } + if !ack.Ok || ack.Seq != 42 { + t.Fatalf("ack = %+v", ack) + } + if gotToken != "test-token" || gotHost != "h-001" || gotContentType != contentTypeProtobuf { + t.Fatalf("headers: token=%q host=%q ct=%q", gotToken, gotHost, gotContentType) + } +} + +func TestReporterSendMetricsRetriesOnHTTPError(t *testing.T) { + var calls int32 + ts := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + if atomic.AddInt32(&calls, 1) < 3 { + http.Error(w, "temporary", http.StatusInternalServerError) + return + } + w.Header().Set("Content-Type", contentTypeProtobuf) + _, _ = w.Write(ackProto(t, true, 7, "")) + })) + defer ts.Close() + + r := newTestReporter(ts.URL, "test-token") + ack, err := r.SendMetrics(context.Background(), &collectorpb.MetricBatch{HostId: "h-001", Seq: 7}) + if err != nil { + t.Fatalf("SendMetrics: %v", err) + } + if !ack.Ok || ack.Seq != 7 { + t.Fatalf("ack = %+v", ack) + } + if atomic.LoadInt32(&calls) != 3 { + t.Fatalf("calls = %d, want 3", atomic.LoadInt32(&calls)) + } +} + +func TestReporterSendMetricsAckNotOk(t *testing.T) { + ts := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + w.Header().Set("Content-Type", contentTypeProtobuf) + _, _ = w.Write(ackProto(t, false, 0, "storage unavailable")) + })) + defer ts.Close() + + r := newTestReporter(ts.URL, "test-token") + _, err := r.SendMetrics(context.Background(), &collectorpb.MetricBatch{HostId: "h-001", Seq: 1}) + if err == nil { + t.Fatalf("expected error for ack not ok") + } + if !strings.Contains(err.Error(), "storage unavailable") { + t.Fatalf("err = %v", err) + } +} + +func TestReporterSendMetricsExhaustsRetries(t *testing.T) { + var calls int32 + ts := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + atomic.AddInt32(&calls, 1) + http.Error(w, "unavailable", http.StatusServiceUnavailable) + })) + defer ts.Close() + + r := newTestReporter(ts.URL, "test-token") + _, err := r.SendMetrics(context.Background(), &collectorpb.MetricBatch{HostId: "h-001", Seq: 1}) + if err == nil { + t.Fatalf("expected error after retries") + } + if atomic.LoadInt32(&calls) != 3 { + t.Fatalf("calls = %d, want 3", atomic.LoadInt32(&calls)) + } + if !strings.Contains(err.Error(), "failed after 3 attempts") { + t.Fatalf("err = %v", err) + } +} + +func TestReporterSendLogsSuccess(t *testing.T) { + ts := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + body, _ := io.ReadAll(r.Body) + var batch collectorpb.LogBatch + if err := proto.Unmarshal(body, &batch); err != nil { + t.Errorf("unmarshal log batch: %v", err) + } + if batch.HostId != "h-001" || batch.Seq != 9 || len(batch.Entries) != 1 { + t.Errorf("batch = %+v", &batch) + } + w.Header().Set("Content-Type", contentTypeProtobuf) + _, _ = w.Write(ackProto(t, true, 9, "")) + })) + defer ts.Close() + + r := newTestReporter(ts.URL, "test-token") + ack, err := r.SendLogs(context.Background(), &collectorpb.LogBatch{ + HostId: "h-001", + Seq: 9, + Entries: []*collectorpb.LogEntry{{Level: "ERROR", Message: "boom"}}, + }) + if err != nil { + t.Fatalf("SendLogs: %v", err) + } + if !ack.Ok || ack.Seq != 9 { + t.Fatalf("ack = %+v", ack) + } +} + +func TestReporterHeartbeatSuccess(t *testing.T) { + ts := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + w.Header().Set("Content-Type", contentTypeProtobuf) + data, _ := proto.Marshal(&collectorpb.HeartbeatResponse{Ok: true, ServerTime: 1700000000}) + _, _ = w.Write(data) + })) + defer ts.Close() + + r := newTestReporter(ts.URL, "test-token") + hb, err := r.Heartbeat(context.Background(), 0) + if err != nil { + t.Fatalf("Heartbeat: %v", err) + } + if !hb.Ok || hb.ServerTime != 1700000000 { + t.Fatalf("hb = %+v", hb) + } +}