test(e2e): 第5刀b L4 gRPC 全链路(enroll→ReportUsage→统计真入库真读出)

在 HTTP 段基础上接 gRPC mTLS 全链路,真验证「统计准不准」:
- 脚本:生成 Node CA + gRPC server 证书(镜像 deploy.sh)、起 server 带 gRPC、
  nodectl 签 bootstrap token,传 E2E_GRPC_ADDR/CA/token/node 给 driver
- driver:复用 agentd.EnsureEnrolled 真 TCP enroll → enroll 到的 client cert
  建 mTLS → ReportUsage 注入 1GiB/100MiB(dp_uuid=用户) → 轮询 /v1/usage
  断言 bytes_up=1GiB·bytes_down=100MiB 真入库真读回
- 本地全链路绿;go build ./... 不受影响(tag e2e)

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
This commit is contained in:
wangjia
2026-06-25 07:04:26 +08:00
parent b3cbeef7cf
commit fd82d82fe3
2 changed files with 158 additions and 19 deletions
+114
View File
@@ -8,6 +8,9 @@ package e2e
import (
"bytes"
"context"
"crypto/tls"
"crypto/x509"
"encoding/json"
"io"
"net/http"
@@ -15,8 +18,12 @@ import (
"testing"
"time"
"github.com/wangjia/pangolin/server/internal/agentd"
"github.com/wangjia/pangolin/server/internal/auth"
"github.com/wangjia/pangolin/server/internal/db"
agentv1 "github.com/wangjia/pangolin/server/internal/pb/agentv1"
"google.golang.org/grpc"
"google.golang.org/grpc/credentials"
)
const (
@@ -59,8 +66,115 @@ func TestE2ESmoke(t *testing.T) {
assertUnauthorized(t, base, "/v1/me")
t.Log("✓ HTTP 段通过:登录 → /v1/me(dp_uuid/plan/email) → /v1/usage(鉴权)")
// ── gRPC 全链路:enroll → ReportUsage → /v1/usage bytes 断言 ──────
if os.Getenv("E2E_GRPC_ADDR") != "" {
runGRPCChain(t, base, token)
} else {
t.Log("E2E_GRPC_ADDR 未设,跳过 gRPC 全链路段")
}
}
// runGRPCChain 镜像真 agent:enroll(mTLS,bootstrap token)→ 用 enroll 到的 client
// cert 调 ReportUsage 注入用量 → 轮询 HTTP /v1/usage 断言 bytes 真入库真读出。
func runGRPCChain(t *testing.T, base, token string) {
t.Helper()
grpcAddr := os.Getenv("E2E_GRPC_ADDR")
caPEM, err := os.ReadFile(os.Getenv("E2E_CA_CERT"))
if err != nil {
t.Fatalf("读 CA cert: %v", err)
}
ctx, cancel := context.WithTimeout(context.Background(), 20*time.Second)
defer cancel()
// 1. enroll(复用生产 agent 的 EnsureEnrolled,真 TCP enroll dialer)
cfg := agentd.Config{
BootstrapToken: os.Getenv("E2E_BOOTSTRAP_TOKEN"),
StateDir: t.TempDir(),
ServerName: "localhost",
AgentVersion: "e2e-smoke",
}
nodeUUID, err := agentd.EnsureEnrolled(ctx, cfg, enrollDialer(grpcAddr, caPEM))
if err != nil {
t.Fatalf("enroll: %v", err)
}
if nodeUUID == "" {
t.Fatal("enroll 返回空 node uuid")
}
t.Logf("✓ enroll 成功 node=%s", nodeUUID)
// 2. 用 enroll 到的 client cert 建 mTLS 连接 → ReportUsage
cert, err := tls.LoadX509KeyPair(cfg.CertPath(), cfg.KeyPath())
if err != nil {
t.Fatalf("加载 node 证书: %v", err)
}
pool := x509.NewCertPool()
if !pool.AppendCertsFromPEM(caPEM) {
t.Fatal("追加 CA 失败")
}
conn, err := grpc.NewClient(grpcAddr, grpc.WithTransportCredentials(
credentials.NewTLS(&tls.Config{
Certificates: []tls.Certificate{cert},
RootCAs: pool,
ServerName: "localhost",
}),
))
if err != nil {
t.Fatalf("mTLS 拨号: %v", err)
}
defer conn.Close()
const wantUp, wantDown = int64(1) << 30, int64(100) << 20 // 1GiB / 100MiB
now := time.Now().Unix()
if _, err := agentv1.NewAgentServiceClient(conn).ReportUsage(ctx, &agentv1.UsageReport{
NodeUUID: nodeUUID,
WindowStartUnix: now - 60,
WindowEndUnix: now,
Entries: []*agentv1.UsageEntry{{
DpUUID: e2eDPUUID, BytesUp: wantUp, BytesDown: wantDown, SessionMinutes: 1,
}},
}); err != nil {
t.Fatalf("ReportUsage: %v", err)
}
t.Log("✓ ReportUsage 上报 1GiB/100MiB")
// 3. 轮询 /v1/usage 直到用量入库可读(聚合是同步写,留少量重试容错)
var up, down int64
for i := 0; i < 30; i++ {
up, down = sumUsage(getJSON(t, base, "/v1/usage", token))
if up == wantUp && down == wantDown {
t.Logf("✓ /v1/usage 读回 bytes_up=%d bytes_down=%d —— gRPC 全链路通", up, down)
return
}
time.Sleep(150 * time.Millisecond)
}
t.Fatalf("/v1/usage 最终 bytes_up=%d(want %d) bytes_down=%d(want %d)", up, wantUp, down, wantDown)
}
// enrollDialer 返回未认证(无 client cert,bootstrap token 鉴权)的 enroll 连接,
// 仅信任 server CA。enroll 拿到 client cert 后,后续 RPC 才走 mTLS。
func enrollDialer(addr string, caPEM []byte) agentd.EnrollDialer {
return func(_ context.Context) (*grpc.ClientConn, error) {
pool := x509.NewCertPool()
pool.AppendCertsFromPEM(caPEM)
return grpc.NewClient(addr, grpc.WithTransportCredentials(
credentials.NewTLS(&tls.Config{RootCAs: pool, ServerName: "localhost"}),
))
}
}
func sumUsage(m map[string]any) (up, down int64) {
pts, _ := m["points"].([]any)
for _, p := range pts {
pm, _ := p.(map[string]any)
up += toInt64(pm["bytes_up"])
down += toInt64(pm["bytes_down"])
}
return up, down
}
func toInt64(v any) int64 { f, _ := v.(float64); return int64(f) }
// seedUser 直接写 sqlite 造一个 active 用户(复用 argon2id HashPassword),
// 比走注册验证码稳。server 刚起基本空闲,一次性 INSERT 锁风险极低。
func seedUser(t *testing.T, dsn string) {