136 lines
3.8 KiB
Go
Raw Permalink Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

package agent
import (
"bytes"
"context"
"fmt"
"io"
"net/http"
"time"
"google.golang.org/protobuf/proto"
"git.opencomputing.cn/yumoqing/host-metrics-collector/pkg/collectorpb"
)
// Agent 上报使用的 HTTP 头与内容类型,与 gateway 包及测试保持一致。
const (
headerAgentToken = "X-Agent-Token"
headerHostID = "X-Host-Id"
contentTypeProtobuf = "application/x-protobuf"
)
// Reporter 负责通过 HTTP/protobuf 向 Collector Gateway(4318)上报数据。
type Reporter struct {
client *http.Client
gatewayURL string
agentToken string
hostID string
agentVersion string
maxAttempts int
baseBackoff time.Duration
}
// NewReporter 创建上报客户端。
func NewReporter(gatewayURL, agentToken, hostID, agentVersion string) *Reporter {
return &Reporter{
client: &http.Client{Timeout: 10 * time.Second},
gatewayURL: gatewayURL,
agentToken: agentToken,
hostID: hostID,
agentVersion: agentVersion,
maxAttempts: 5,
baseBackoff: time.Second,
}
}
// Heartbeat 发送心跳,返回服务端时间与下发配置。
func (r *Reporter) Heartbeat(ctx context.Context, configVersion int64) (*collectorpb.HeartbeatResponse, error) {
req := &collectorpb.HeartbeatRequest{HostId: r.hostID, AgentVersion: r.agentVersion, ConfigVersion: configVersion}
out := &collectorpb.HeartbeatResponse{}
if err := r.sendWithRetry(ctx, "/v1/agent/heartbeat", req, out); err != nil {
return nil, err
}
return out, nil
}
// SendMetrics 上报指标批次,返回服务端 Ack。
func (r *Reporter) SendMetrics(ctx context.Context, batch *collectorpb.MetricBatch) (*collectorpb.Ack, error) {
out := &collectorpb.Ack{}
if err := r.sendWithRetry(ctx, "/v1/agent/metrics", batch, out); err != nil {
return nil, err
}
return out, nil
}
// SendLogs 上报日志批次,返回服务端 Ack。
func (r *Reporter) SendLogs(ctx context.Context, batch *collectorpb.LogBatch) (*collectorpb.Ack, error) {
out := &collectorpb.Ack{}
if err := r.sendWithRetry(ctx, "/v1/agent/logs", batch, out); err != nil {
return nil, err
}
return out, nil
}
func (r *Reporter) sendWithRetry(ctx context.Context, path string, req proto.Message, out proto.Message) error {
var lastErr error
for attempt := 0; attempt < r.maxAttempts; attempt++ {
if attempt > 0 {
backoff := r.baseBackoff << (attempt - 1)
if backoff > 60*time.Second {
backoff = 60 * time.Second
}
select {
case <-ctx.Done():
return ctx.Err()
case <-time.After(backoff):
}
}
err := r.post(ctx, path, req, out)
if err == nil {
if ack, ok := out.(*collectorpb.Ack); ok && !ack.Ok {
lastErr = fmt.Errorf("gateway ack error: %s", ack.Error)
continue
}
return nil
}
lastErr = err
}
return fmt.Errorf("report %s failed after %d attempts: %w", path, r.maxAttempts, lastErr)
}
func (r *Reporter) post(ctx context.Context, path string, req proto.Message, out proto.Message) error {
body, err := proto.Marshal(req)
if err != nil {
return fmt.Errorf("marshal request: %w", err)
}
httpReq, err := http.NewRequestWithContext(ctx, http.MethodPost, r.gatewayURL+path, bytes.NewReader(body))
if err != nil {
return fmt.Errorf("new request: %w", err)
}
httpReq.Header.Set("Content-Type", contentTypeProtobuf)
httpReq.Header.Set(headerAgentToken, r.agentToken)
httpReq.Header.Set(headerHostID, r.hostID)
resp, err := r.client.Do(httpReq)
if err != nil {
return fmt.Errorf("do request: %w", err)
}
defer resp.Body.Close()
data, err := io.ReadAll(io.LimitReader(resp.Body, 4<<20))
if err != nil {
return fmt.Errorf("read response: %w", err)
}
if resp.StatusCode != http.StatusOK {
return fmt.Errorf("gateway returned %d: %s", resp.StatusCode, string(data))
}
if out != nil && len(data) > 0 {
if err := proto.Unmarshal(data, out); err != nil {
return fmt.Errorf("unmarshal response: %w", err)
}
}
return nil
}