136 lines
3.8 KiB
Go
136 lines
3.8 KiB
Go
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
|
||
}
|