feat(gateway): 网关健康状态——后台定时探活只保留最近一次结果并在列表展示
- 迁移 1043:gateway 表加 health_status/last_check_reason/last_checked_at - Store.RecordGatewayCheck 覆盖写最近一次探活结果,Create/Update/List/Get 全量返回健康字段 - api.RecordGatewayHealthChecks 并发探 /healthz:连接失败或 5xx 判离线,其余判在线(401 等异常写入原因供排查) - workers 每 30 秒一轮,app 启动接线 - 网关管理列表新增健康状态列:在线/离线/未检查 Tag,原因与检查时间收进 Tooltip
This commit is contained in:
@@ -0,0 +1,72 @@
|
||||
package api
|
||||
|
||||
import (
|
||||
"context"
|
||||
"net/http"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
hub "git.ipao.vip/rogee/creator-hub/internal/environment"
|
||||
"github.com/sirupsen/logrus"
|
||||
)
|
||||
|
||||
const gatewayHealthProbeTimeout = 5 * time.Second
|
||||
|
||||
// gatewayHealthStore 是网关健康检查所需的存储能力;生产实现为 *hub.Store,测试使用内存桩。
|
||||
type gatewayHealthStore interface {
|
||||
ListGateways(ctx context.Context) ([]hub.Gateway, error)
|
||||
RecordGatewayCheck(ctx context.Context, name, status, reason string) (hub.Gateway, error)
|
||||
}
|
||||
|
||||
// classifyGatewayHealth 将一次 /healthz 探活结果归一为持久化状态与原因。
|
||||
// 连接失败或 5xx 判为 unhealthy;其余(含 401 令牌不匹配)说明进程可达,判为 healthy
|
||||
// 并把异常状态写进 reason 供排查。
|
||||
func classifyGatewayHealth(status int, callErr error) (string, string) {
|
||||
if callErr != nil {
|
||||
reason := callErr.Error()
|
||||
if len(reason) > 300 {
|
||||
reason = reason[:300]
|
||||
}
|
||||
return "unhealthy", reason
|
||||
}
|
||||
if status >= 500 {
|
||||
return "unhealthy", "gateway returned status " + http.StatusText(status)
|
||||
}
|
||||
if status >= 400 {
|
||||
return "healthy", "gateway responded with status " + http.StatusText(status)
|
||||
}
|
||||
return "healthy", ""
|
||||
}
|
||||
|
||||
// RecordGatewayHealthChecks 执行一轮网关健康探活并落库;供 workers 定时任务与测试调用。
|
||||
func RecordGatewayHealthChecks(ctx context.Context, store *hub.Store) error {
|
||||
return recordGatewayHealthChecks(ctx, store)
|
||||
}
|
||||
|
||||
// recordGatewayHealthChecks 对全部注册网关并发探活一次,每个网关只覆盖写入最近一次结果。
|
||||
// 单个网关落库失败不阻断其余网关,仅记日志(下一轮探活会重试)。
|
||||
func recordGatewayHealthChecks(ctx context.Context, store gatewayHealthStore) error {
|
||||
gateways, err := store.ListGateways(ctx)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
var wg sync.WaitGroup
|
||||
for _, gateway := range gateways {
|
||||
wg.Add(1)
|
||||
go func() {
|
||||
defer wg.Done()
|
||||
status, _, callErr := gatewayCall(ctx, gateway, http.MethodGet, "/healthz", nil, gatewayHealthProbeTimeout)
|
||||
healthStatus, reason := classifyGatewayHealth(status, callErr)
|
||||
if _, recordErr := store.RecordGatewayCheck(ctx, gateway.Name, healthStatus, reason); recordErr != nil {
|
||||
logrus.WithFields(logrus.Fields{
|
||||
"service": "control-plane",
|
||||
"event_type": "gateway_health_check",
|
||||
"gateway": gateway.Name,
|
||||
"probed": healthStatus,
|
||||
}).WithError(recordErr).Warn("gateway health check result was not persisted")
|
||||
}
|
||||
}()
|
||||
}
|
||||
wg.Wait()
|
||||
return nil
|
||||
}
|
||||
@@ -0,0 +1,122 @@
|
||||
package api
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
"sync"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
hub "git.ipao.vip/rogee/creator-hub/internal/environment"
|
||||
)
|
||||
|
||||
type recordedGatewayCheck struct {
|
||||
name string
|
||||
status string
|
||||
reason string
|
||||
}
|
||||
|
||||
// healthMemoryStore 在 memoryStore 之上记录 RecordGatewayCheck 调用,供健康检查断言。
|
||||
type healthMemoryStore struct {
|
||||
*memoryStore
|
||||
mu sync.Mutex
|
||||
checks []recordedGatewayCheck
|
||||
}
|
||||
|
||||
func (s *healthMemoryStore) RecordGatewayCheck(_ context.Context, name, status, reason string) (hub.Gateway, error) {
|
||||
s.mu.Lock()
|
||||
defer s.mu.Unlock()
|
||||
s.checks = append(s.checks, recordedGatewayCheck{name: name, status: status, reason: reason})
|
||||
return hub.Gateway{Name: name, HealthStatus: status, LastCheckReason: reason}, nil
|
||||
}
|
||||
|
||||
func TestClassifyGatewayHealth(t *testing.T) {
|
||||
for _, test := range []struct {
|
||||
name string
|
||||
status int
|
||||
callErr error
|
||||
wantStatus string
|
||||
wantReasonHad bool // 是否应携带非空 reason
|
||||
}{
|
||||
{name: "204 healthy without reason", status: http.StatusNoContent, wantStatus: "healthy", wantReasonHad: false},
|
||||
{name: "200 healthy without reason", status: http.StatusOK, wantStatus: "healthy", wantReasonHad: false},
|
||||
{name: "401 reachable but unauthorized", status: http.StatusUnauthorized, wantStatus: "healthy", wantReasonHad: true},
|
||||
{name: "404 reachable but misrouted", status: http.StatusNotFound, wantStatus: "healthy", wantReasonHad: true},
|
||||
{name: "503 unhealthy", status: http.StatusServiceUnavailable, wantStatus: "unhealthy", wantReasonHad: true},
|
||||
{name: "connection error unhealthy", callErr: errors.New(`Get "http://gw:8081/healthz": dial tcp: connection refused`), wantStatus: "unhealthy", wantReasonHad: true},
|
||||
} {
|
||||
t.Run(test.name, func(t *testing.T) {
|
||||
status, reason := classifyGatewayHealth(test.status, test.callErr)
|
||||
if status != test.wantStatus {
|
||||
t.Fatalf("classify status = %q, want %q (reason=%q)", status, test.wantStatus, reason)
|
||||
}
|
||||
if test.wantReasonHad != (reason != "") {
|
||||
t.Fatalf("classify reason = %q, wantReasonHad=%v", reason, test.wantReasonHad)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func TestClassifyGatewayHealthTruncatesReason(t *testing.T) {
|
||||
long := errors.New(string(make([]byte, 1024)))
|
||||
_, reason := classifyGatewayHealth(0, long)
|
||||
if len(reason) != 300 {
|
||||
t.Fatalf("reason length = %d, want 300", len(reason))
|
||||
}
|
||||
}
|
||||
|
||||
func TestRecordGatewayHealthChecksPersistsLatestResult(t *testing.T) {
|
||||
reachable := &fakeGateway{token: "token-a"}
|
||||
reachableServer := httptest.NewServer(reachable.handler(t))
|
||||
t.Cleanup(reachableServer.Close)
|
||||
|
||||
closedServer := httptest.NewServer(http.HandlerFunc(func(http.ResponseWriter, *http.Request) {}))
|
||||
closedServer.Close() // 端口已释放,保证连接失败
|
||||
|
||||
store := &healthMemoryStore{memoryStore: newMemoryStore()}
|
||||
store.gateways["gw-live"] = hub.Gateway{Name: "gw-live", Endpoint: reachableServer.URL, Token: reachable.token}
|
||||
store.gateways["gw-down"] = hub.Gateway{Name: "gw-down", Endpoint: closedServer.URL, Token: "unused-token"}
|
||||
|
||||
if err := recordGatewayHealthChecks(context.Background(), store); err != nil {
|
||||
t.Fatalf("record gateway health checks returned error: %v", err)
|
||||
}
|
||||
|
||||
byName := map[string]recordedGatewayCheck{}
|
||||
for _, check := range store.checks {
|
||||
byName[check.name] = check
|
||||
}
|
||||
if len(store.checks) != 2 {
|
||||
t.Fatalf("expected one check per gateway, got %#v", store.checks)
|
||||
}
|
||||
if byName["gw-live"].status != "healthy" {
|
||||
t.Fatalf("reachable gateway status = %q, want healthy", byName["gw-live"].status)
|
||||
}
|
||||
if byName["gw-live"].reason != "" {
|
||||
t.Fatalf("reachable gateway reason = %q, want empty", byName["gw-live"].reason)
|
||||
}
|
||||
if byName["gw-down"].status != "unhealthy" || byName["gw-down"].reason == "" {
|
||||
t.Fatalf("unreachable gateway = %#v, want unhealthy with reason", byName["gw-down"])
|
||||
}
|
||||
if requests := reachable.recorded(); len(requests) != 1 || requests[0].path != "/healthz" {
|
||||
t.Fatalf("reachable gateway probes = %#v, want exactly one /healthz", requests)
|
||||
}
|
||||
}
|
||||
|
||||
func TestRecordGatewayHealthChecksEmptyListIsNoop(t *testing.T) {
|
||||
store := &healthMemoryStore{memoryStore: newMemoryStore()}
|
||||
if err := recordGatewayHealthChecks(context.Background(), store); err != nil {
|
||||
t.Fatalf("empty list returned error: %v", err)
|
||||
}
|
||||
if len(store.checks) != 0 {
|
||||
t.Fatalf("empty list recorded checks: %#v", store.checks)
|
||||
}
|
||||
}
|
||||
|
||||
// 编译期约束:探活超时常量必须明显小于调度间隔,避免探活互相堆叠。
|
||||
func TestGatewayHealthProbeTimeoutBounded(t *testing.T) {
|
||||
if gatewayHealthProbeTimeout > 10*time.Second {
|
||||
t.Fatalf("probe timeout %v must stay below scheduling cadence", gatewayHealthProbeTimeout)
|
||||
}
|
||||
}
|
||||
@@ -573,6 +573,9 @@ func (probe *sequenceExitProbe) Check(context.Context, hub.NetworkExitAccess) (h
|
||||
func (g *fakeGateway) handler(t *testing.T) http.Handler {
|
||||
return http.HandlerFunc(func(response http.ResponseWriter, request *http.Request) {
|
||||
if request.Method == http.MethodGet && request.URL.Path == "/healthz" {
|
||||
g.mu.Lock()
|
||||
g.requests = append(g.requests, recordedRequest{method: request.Method, path: request.URL.Path})
|
||||
g.mu.Unlock()
|
||||
response.WriteHeader(http.StatusNoContent)
|
||||
return
|
||||
}
|
||||
@@ -829,12 +832,13 @@ func TestUpdateGatewayRenamesAndPreservesReferences(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestListGatewaysReturnsStoredRecordsWithoutProbing(t *testing.T) {
|
||||
func TestListGatewaysReturnsStoredHealthWithoutProbing(t *testing.T) {
|
||||
store := newMemoryStore()
|
||||
gateway := &fakeGateway{token: "unit-test-gateway-token"}
|
||||
server := httptest.NewServer(gateway.handler(t))
|
||||
t.Cleanup(server.Close)
|
||||
stored := hub.Gateway{Name: "gw-1", Endpoint: server.URL, Token: gateway.token}
|
||||
stored := hub.Gateway{Name: "gw-1", Endpoint: server.URL, Token: gateway.token,
|
||||
HealthStatus: "healthy", LastCheckReason: "", LastCheckedAt: nil}
|
||||
store.gateways[stored.Name] = stored
|
||||
app := fiber.New()
|
||||
registerHubWithNetwork(app, store, fakeExitProbe{}, func(hub.NetworkExitAccess) (string, error) { return "", nil })
|
||||
@@ -847,8 +851,8 @@ func TestListGatewaysReturnsStoredRecordsWithoutProbing(t *testing.T) {
|
||||
if err := json.Unmarshal(response.Body.Bytes(), &gateways); err != nil || len(gateways) != 1 || gateways[0] != stored {
|
||||
t.Fatalf("stored gateway was not returned unchanged: gateways=%#v err=%v", gateways, err)
|
||||
}
|
||||
if strings.Contains(response.Body.String(), "connectivity") || strings.Contains(response.Body.String(), "health") {
|
||||
t.Fatalf("gateway list exposed live status fields: %s", response.Body.String())
|
||||
if !strings.Contains(response.Body.String(), `"health_status":"healthy"`) {
|
||||
t.Fatalf("gateway list should return persisted health status: %s", response.Body.String())
|
||||
}
|
||||
if requests := gateway.recorded(); len(requests) != 0 {
|
||||
t.Fatalf("gateway list performed live probes: %#v", requests)
|
||||
|
||||
@@ -97,13 +97,21 @@ func newCommand() *cobra.Command {
|
||||
defer close(creatorScheduleDone)
|
||||
workers.RunCreatorScheduler(creatorScheduleContext, creatorStore, phaseAStore, hubStore)
|
||||
}()
|
||||
gatewayHealthContext, stopGatewayHealth := context.WithCancel(command.Context())
|
||||
gatewayHealthDone := make(chan struct{})
|
||||
go func() {
|
||||
defer close(gatewayHealthDone)
|
||||
workers.RunGatewayHealthCheck(gatewayHealthContext, hubStore)
|
||||
}()
|
||||
listenErr := newHandler(cfg.webDir, cfg.username, cfg.password, phaseAStore, hubStore, credentials, creatorStore).Listen(cfg.listenAddr, fiber.ListenConfig{
|
||||
GracefulContext: command.Context(),
|
||||
DisableStartupMessage: true,
|
||||
})
|
||||
stopCreatorScheduler()
|
||||
stopGatewayHealth()
|
||||
stopHeartbeat()
|
||||
<-creatorScheduleDone
|
||||
<-gatewayHealthDone
|
||||
<-heartbeatDone
|
||||
return listenErr
|
||||
},
|
||||
|
||||
@@ -29,3 +29,19 @@ func RunRuntimeLeaseHeartbeat(ctx context.Context, store *environment.Store) {
|
||||
func RunCreatorScheduler(ctx context.Context, store *creatorDomain.Store, accountStore *account.Store, environmentStore *environment.Store) {
|
||||
api.RunCreatorScheduler(ctx, store, accountStore, environmentStore)
|
||||
}
|
||||
|
||||
// RunGatewayHealthCheck 每 30 秒对全部注册网关探活一次,只保留最近一次结果。
|
||||
func RunGatewayHealthCheck(ctx context.Context, store *environment.Store) {
|
||||
ticker := time.NewTicker(30 * time.Second)
|
||||
defer ticker.Stop()
|
||||
for {
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
return
|
||||
case <-ticker.C:
|
||||
if err := api.RecordGatewayHealthChecks(ctx, store); err != nil && ctx.Err() == nil {
|
||||
logrus.WithField("service", "control-plane").WithError(err).Warn("gateway health check failed")
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -70,6 +70,6 @@ func TestMigration043ConsolidatedSchemaShape(t *testing.T) {
|
||||
// audit_event 任务列已删
|
||||
assertDatabaseCount(t, db, `SELECT count(*) FROM information_schema.columns WHERE table_schema = current_schema()
|
||||
AND table_name = 'audit_event' AND column_name IN ('confirmation_id','confirmation_version','attempt_id','task_id','runtime_instance_id')`, 0)
|
||||
// 统一登记表:1-38(除 36)、1017-1042、43 全部登记
|
||||
assertDatabaseCount(t, db, `SELECT count(*) FROM schema_migration`, 49)
|
||||
// 统一登记表:1-38(除 36)、1017-1043、43 全部登记
|
||||
assertDatabaseCount(t, db, `SELECT count(*) FROM schema_migration`, 50)
|
||||
}
|
||||
|
||||
@@ -30,7 +30,7 @@ func TestUnifiedAccountMigration(t *testing.T) {
|
||||
t.Fatal(err)
|
||||
}
|
||||
defer db.Close()
|
||||
assertDatabaseCount(t, db, `SELECT count(*) FROM schema_migration`, 49)
|
||||
assertDatabaseCount(t, db, `SELECT count(*) FROM schema_migration`, 50)
|
||||
assertDatabaseCount(t, db, `SELECT count(*) FROM information_schema.tables WHERE table_schema = current_schema()
|
||||
AND table_name IN ('social_account', 'browser_env', 'network_exit', 'environment_binding')`, 3)
|
||||
assertDatabaseCount(t, db, `SELECT count(*) FROM information_schema.tables WHERE table_schema = current_schema() AND table_name = 'browser_image'`, 0)
|
||||
@@ -42,7 +42,7 @@ func TestUnifiedAccountMigration(t *testing.T) {
|
||||
|
||||
store = openFullyMigratedHub(t, ctx, testURL)
|
||||
store.Close()
|
||||
assertDatabaseCount(t, db, `SELECT count(*) FROM schema_migration`, 49)
|
||||
assertDatabaseCount(t, db, `SELECT count(*) FROM schema_migration`, 50)
|
||||
})
|
||||
|
||||
t.Run("legacy migration 013 without account secrets is repaired forward", func(t *testing.T) {
|
||||
|
||||
@@ -0,0 +1,5 @@
|
||||
-- 网关健康状态:后台定时探活只保留最近一次结果(不保留历史),列表接口直接读库展示。
|
||||
ALTER TABLE gateway ADD COLUMN health_status text NOT NULL DEFAULT 'unchecked'
|
||||
CHECK (health_status IN ('unchecked', 'healthy', 'unhealthy'));
|
||||
ALTER TABLE gateway ADD COLUMN last_check_reason text NOT NULL DEFAULT '';
|
||||
ALTER TABLE gateway ADD COLUMN last_checked_at timestamptz;
|
||||
@@ -167,6 +167,9 @@ var migration1041 string
|
||||
//go:embed migrations/1042_work_cover_cache.sql
|
||||
var migration1042 string
|
||||
|
||||
//go:embed migrations/1043_gateway_health.sql
|
||||
var migration1043 string
|
||||
|
||||
var (
|
||||
ErrConflict = errors.New("resource conflicts with existing state")
|
||||
ErrInvalid = errors.New("invalid environment input")
|
||||
@@ -188,12 +191,16 @@ type Store struct {
|
||||
}
|
||||
|
||||
// Gateway 是平台注册的 native browser gateway 节点;Token 由平台生成,明文存储供页面复制(开发阶段约定)。
|
||||
// 健康字段由控制面后台定时探活写入,只保留最近一次结果,不保留历史。
|
||||
type Gateway struct {
|
||||
Name string `json:"name"`
|
||||
Endpoint string `json:"endpoint"`
|
||||
Token string `json:"token"`
|
||||
CreatedAt time.Time `json:"created_at"`
|
||||
UpdatedAt time.Time `json:"updated_at"`
|
||||
Name string `json:"name"`
|
||||
Endpoint string `json:"endpoint"`
|
||||
Token string `json:"token"`
|
||||
HealthStatus string `json:"health_status"`
|
||||
LastCheckReason string `json:"last_check_reason"`
|
||||
LastCheckedAt *time.Time `json:"last_checked_at"`
|
||||
CreatedAt time.Time `json:"created_at"`
|
||||
UpdatedAt time.Time `json:"updated_at"`
|
||||
}
|
||||
|
||||
// Env 是一个浏览器环境;浏览器安装和默认运行时由 gateway 宿主机配置。
|
||||
@@ -302,7 +309,7 @@ func (s *Store) migrate(ctx context.Context) error {
|
||||
{1023, migration1023}, {1024, migration1024}, {1025, migration1025}, {1026, migration1026}, {1027, migration1027}, {1028, migration1028},
|
||||
{1029, migration1029}, {1030, migration1030}, {1031, migration1031}, {1032, migration1032}, {1033, migration1033}, {1034, migration1034},
|
||||
{1035, migration1035}, {1036, migration1036}, {1037, migration1037}, {1038, migration1038}, {1039, migration1039}, {1040, migration1040},
|
||||
{1041, migration1041}, {1042, migration1042},
|
||||
{1041, migration1041}, {1042, migration1042}, {1043, migration1043},
|
||||
{43, migration043}, {44, migration044}} {
|
||||
var applied bool
|
||||
if err := tx.QueryRowContext(ctx, `SELECT EXISTS (SELECT 1 FROM schema_migration WHERE version = $1)`, migration.version).Scan(&applied); err != nil {
|
||||
@@ -324,6 +331,19 @@ func (s *Store) migrate(ctx context.Context) error {
|
||||
return nil
|
||||
}
|
||||
|
||||
func scanGateway(row rowScanner) (Gateway, error) {
|
||||
var gateway Gateway
|
||||
var checked sql.NullTime
|
||||
if err := row.Scan(&gateway.Name, &gateway.Endpoint, &gateway.Token, &gateway.HealthStatus,
|
||||
&gateway.LastCheckReason, &checked, &gateway.CreatedAt, &gateway.UpdatedAt); err != nil {
|
||||
return Gateway{}, rowError(err)
|
||||
}
|
||||
if checked.Valid {
|
||||
gateway.LastCheckedAt = &checked.Time
|
||||
}
|
||||
return gateway, nil
|
||||
}
|
||||
|
||||
func (s *Store) CreateGateway(ctx context.Context, name, endpoint, token string) (Gateway, error) {
|
||||
name, endpoint, token = strings.TrimSpace(name), strings.TrimSpace(endpoint), strings.TrimSpace(token)
|
||||
if !gatewayNamePattern.MatchString(name) || !validHTTPURL(endpoint) {
|
||||
@@ -337,15 +357,13 @@ func (s *Store) CreateGateway(ctx context.Context, name, endpoint, token string)
|
||||
} else {
|
||||
token = newToken()
|
||||
}
|
||||
gateway := Gateway{Name: name, Endpoint: endpoint, Token: token}
|
||||
err := s.db.QueryRowContext(ctx, `
|
||||
result, err := scanGateway(s.db.QueryRowContext(ctx, `
|
||||
INSERT INTO gateway (name, endpoint, token) VALUES ($1, $2, $3)
|
||||
RETURNING created_at, updated_at`, name, endpoint, gateway.Token).
|
||||
Scan(&gateway.CreatedAt, &gateway.UpdatedAt)
|
||||
RETURNING name, endpoint, token, health_status, last_check_reason, last_checked_at, created_at, updated_at`, name, endpoint, token))
|
||||
if err != nil {
|
||||
return Gateway{}, publicDatabaseError(err)
|
||||
return Gateway{}, err
|
||||
}
|
||||
return gateway, nil
|
||||
return result, nil
|
||||
}
|
||||
|
||||
// UpdateGateway 修改网关名称、Endpoint 和令牌。名称变更由数据库外键 ON UPDATE CASCADE
|
||||
@@ -357,31 +375,26 @@ func (s *Store) UpdateGateway(ctx context.Context, currentName, name, endpoint,
|
||||
(token != "" && !tokenPattern.MatchString(token)) {
|
||||
return Gateway{}, ErrInvalid
|
||||
}
|
||||
var gateway Gateway
|
||||
err := s.db.QueryRowContext(ctx, `
|
||||
return scanGateway(s.db.QueryRowContext(ctx, `
|
||||
UPDATE gateway SET name = $1, endpoint = $2,
|
||||
token = CASE WHEN $3 = '' THEN token ELSE $3 END, updated_at = now()
|
||||
WHERE name = $4
|
||||
RETURNING name, endpoint, token, created_at, updated_at`, name, endpoint, token, currentName).
|
||||
Scan(&gateway.Name, &gateway.Endpoint, &gateway.Token, &gateway.CreatedAt, &gateway.UpdatedAt)
|
||||
if err != nil {
|
||||
return Gateway{}, rowError(err)
|
||||
}
|
||||
return gateway, nil
|
||||
RETURNING name, endpoint, token, health_status, last_check_reason, last_checked_at, created_at, updated_at`, name, endpoint, token, currentName))
|
||||
}
|
||||
|
||||
func (s *Store) ListGateways(ctx context.Context) ([]Gateway, error) {
|
||||
rows, err := s.db.QueryContext(ctx, `
|
||||
SELECT name, endpoint, token, created_at, updated_at FROM gateway ORDER BY created_at, name`)
|
||||
SELECT name, endpoint, token, health_status, last_check_reason, last_checked_at, created_at, updated_at
|
||||
FROM gateway ORDER BY created_at, name`)
|
||||
if err != nil {
|
||||
return nil, errors.New("read gateways")
|
||||
}
|
||||
defer rows.Close()
|
||||
gateways := []Gateway{}
|
||||
for rows.Next() {
|
||||
var gateway Gateway
|
||||
if err := rows.Scan(&gateway.Name, &gateway.Endpoint, &gateway.Token, &gateway.CreatedAt, &gateway.UpdatedAt); err != nil {
|
||||
return nil, errors.New("decode gateway")
|
||||
gateway, err := scanGateway(rows)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
gateways = append(gateways, gateway)
|
||||
}
|
||||
@@ -389,17 +402,23 @@ func (s *Store) ListGateways(ctx context.Context) ([]Gateway, error) {
|
||||
}
|
||||
|
||||
func (s *Store) GetGateway(ctx context.Context, name string) (Gateway, error) {
|
||||
var gateway Gateway
|
||||
if !gatewayNamePattern.MatchString(name) {
|
||||
return gateway, ErrInvalid
|
||||
return Gateway{}, ErrInvalid
|
||||
}
|
||||
err := s.db.QueryRowContext(ctx, `
|
||||
SELECT name, endpoint, token, created_at, updated_at FROM gateway WHERE name = $1`, name).
|
||||
Scan(&gateway.Name, &gateway.Endpoint, &gateway.Token, &gateway.CreatedAt, &gateway.UpdatedAt)
|
||||
if err != nil {
|
||||
return Gateway{}, rowError(err)
|
||||
return scanGateway(s.db.QueryRowContext(ctx, `
|
||||
SELECT name, endpoint, token, health_status, last_check_reason, last_checked_at, created_at, updated_at
|
||||
FROM gateway WHERE name = $1`, name))
|
||||
}
|
||||
|
||||
// RecordGatewayCheck 覆盖写入最近一次探活结果;不保留历史。status 只允许 healthy/unhealthy。
|
||||
func (s *Store) RecordGatewayCheck(ctx context.Context, name, status, reason string) (Gateway, error) {
|
||||
if !gatewayNamePattern.MatchString(name) || (status != "healthy" && status != "unhealthy") {
|
||||
return Gateway{}, ErrInvalid
|
||||
}
|
||||
return gateway, nil
|
||||
return scanGateway(s.db.QueryRowContext(ctx, `
|
||||
UPDATE gateway SET health_status = $2, last_check_reason = $3, last_checked_at = now(), updated_at = now()
|
||||
WHERE name = $1
|
||||
RETURNING name, endpoint, token, health_status, last_check_reason, last_checked_at, created_at, updated_at`, name, status, reason))
|
||||
}
|
||||
|
||||
func (s *Store) DeleteGateway(ctx context.Context, name string) error {
|
||||
|
||||
@@ -808,3 +808,62 @@ func TestNetworkExitLifecycleStoreOperations(t *testing.T) {
|
||||
t.Fatalf("list exits: exits=%+v err=%v", exits, err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestGatewayHealthCheckPersistence(t *testing.T) {
|
||||
databaseURL := os.Getenv("CREATORHUB_POSTGRES_TEST_URL")
|
||||
if databaseURL == "" {
|
||||
t.Skip("set CREATORHUB_POSTGRES_TEST_URL to run PostgreSQL integration coverage")
|
||||
}
|
||||
ctx := context.Background()
|
||||
store := openFullyMigratedHub(t, ctx, isolatedDatabaseURL(t, databaseURL))
|
||||
t.Cleanup(func() { _ = store.Close() })
|
||||
|
||||
gateway, err := store.CreateGateway(ctx, "gw-health", "http://127.0.0.1:8081", "")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if gateway.HealthStatus != "unchecked" || gateway.LastCheckReason != "" || gateway.LastCheckedAt != nil {
|
||||
t.Fatalf("new gateway must start unchecked: %#v", gateway)
|
||||
}
|
||||
|
||||
checked, err := store.RecordGatewayCheck(ctx, "gw-health", "healthy", "")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if checked.HealthStatus != "healthy" || checked.LastCheckReason != "" || checked.LastCheckedAt == nil {
|
||||
t.Fatalf("healthy check was not persisted: %#v", checked)
|
||||
}
|
||||
|
||||
failed, err := store.RecordGatewayCheck(ctx, "gw-health", "unhealthy", `Get "http://127.0.0.1:8081/healthz": connection refused`)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if failed.HealthStatus != "unhealthy" || failed.LastCheckReason == "" || failed.LastCheckedAt == nil {
|
||||
t.Fatalf("unhealthy check was not persisted: %#v", failed)
|
||||
}
|
||||
if !failed.LastCheckedAt.After(*checked.LastCheckedAt) && failed.LastCheckedAt.Equal(*checked.LastCheckedAt) {
|
||||
t.Fatalf("last_checked_at must advance: %v", failed.LastCheckedAt)
|
||||
}
|
||||
|
||||
reread, err := store.GetGateway(ctx, "gw-health")
|
||||
if err != nil || reread.HealthStatus != "unhealthy" || reread.LastCheckReason != failed.LastCheckReason {
|
||||
t.Fatalf("health fields must round trip through get/list: %#v err=%v", reread, err)
|
||||
}
|
||||
listed, err := store.ListGateways(ctx)
|
||||
if err != nil || len(listed) != 1 || listed[0].HealthStatus != "unhealthy" {
|
||||
t.Fatalf("health fields must round trip through list: list=%#v err=%v", listed, err)
|
||||
}
|
||||
|
||||
if _, err := store.RecordGatewayCheck(ctx, "gw-health", "degraded", ""); !errors.Is(err, ErrInvalid) {
|
||||
t.Fatalf("unknown status must be rejected, got %v", err)
|
||||
}
|
||||
if _, err := store.RecordGatewayCheck(ctx, "missing", "healthy", ""); !errors.Is(err, ErrNotFound) {
|
||||
t.Fatalf("missing gateway must surface not found, got %v", err)
|
||||
}
|
||||
|
||||
// 改名/改地址不触碰健康字段(由探活负责覆盖)。
|
||||
renamed, err := store.UpdateGateway(ctx, "gw-health", "gw-health-2", "http://127.0.0.2:8081", "")
|
||||
if err != nil || renamed.HealthStatus != "unhealthy" || renamed.LastCheckReason == "" {
|
||||
t.Fatalf("update must preserve health fields: %#v err=%v", renamed, err)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1,15 +1,30 @@
|
||||
// 网关管理:语义对齐 web.archived GatewaysPage.jsx(注册/编辑 Modal、令牌显隐复制、删除)。
|
||||
import { useCallback, useEffect, useState } from 'react';
|
||||
import { Alert, App, Button, Card, Flex, Form, Input, Modal, Popconfirm, Space, Table, Typography } from 'antd';
|
||||
import { Alert, App, Button, Card, Flex, Form, Input, Modal, Popconfirm, Space, Table, Tag, Tooltip, Typography } from 'antd';
|
||||
import { CopyOutlined, EyeInvisibleOutlined, EyeOutlined, PlusOutlined } from '@ant-design/icons';
|
||||
import type { ColumnsType } from 'antd/es/table';
|
||||
import { create, getList, remove, update } from '@/services/api';
|
||||
import { conflictMessage } from '@/utils/helpers';
|
||||
import { conflictMessage, dateTime } from '@/utils/helpers';
|
||||
import { fixedLeft, fixedRight, tablePagination, tableScroll, useOverflowGrid } from '@/utils/table';
|
||||
|
||||
const namePattern = /^[A-Za-z0-9][A-Za-z0-9._-]{0,63}$/;
|
||||
const tokenPattern = /^[A-Za-z0-9][A-Za-z0-9._-]{15,127}$/;
|
||||
|
||||
// 健康状态列:只展示后台定时探活的最近一次结果,不保留历史。
|
||||
function HealthCell({ status, reason, checkedAt }: { status?: string; reason?: string; checkedAt?: string }) {
|
||||
const label = status === 'healthy' ? '在线' : status === 'unhealthy' ? '离线' : '未检查';
|
||||
const color = status === 'healthy' ? 'success' : status === 'unhealthy' ? 'error' : 'default';
|
||||
const details = [
|
||||
checkedAt ? `最近检查:${dateTime(checkedAt)}` : null,
|
||||
reason || null,
|
||||
].filter(Boolean).join(';');
|
||||
return (
|
||||
<Tooltip title={details || '尚未执行健康检查'}>
|
||||
<Tag color={color}>{label}</Tag>
|
||||
</Tooltip>
|
||||
);
|
||||
}
|
||||
|
||||
function GatewayFormModal({
|
||||
open,
|
||||
initial,
|
||||
@@ -146,6 +161,12 @@ export default function Page() {
|
||||
const columns: ColumnsType<any> = [
|
||||
{ title: '名称', dataIndex: 'name', ...fixedLeft<any>({}), render: (value: string) => <Typography.Text strong>{value}</Typography.Text> },
|
||||
{ title: 'Endpoint', dataIndex: 'endpoint', render: (value: string) => <Typography.Text code style={{ fontSize: 12 }} ellipsis>{value}</Typography.Text> },
|
||||
{
|
||||
title: '健康状态',
|
||||
dataIndex: 'health_status',
|
||||
width: 110,
|
||||
render: (_, gateway) => <HealthCell status={gateway.health_status} reason={gateway.last_check_reason} checkedAt={gateway.last_checked_at} />,
|
||||
},
|
||||
{
|
||||
title: '令牌',
|
||||
dataIndex: 'token',
|
||||
|
||||
@@ -0,0 +1,40 @@
|
||||
const assert = require('node:assert/strict');
|
||||
const { readFileSync } = require('node:fs');
|
||||
const { test } = require('node:test');
|
||||
const ts = require('typescript');
|
||||
|
||||
// 网关健康状态列:后台定时探活只落最近一次结果,列表直接展示,不保留历史。
|
||||
const source = ts.createSourceFile(
|
||||
'gateways-page',
|
||||
readFileSync(`${__dirname}/../src/pages/gateways/index.tsx`, 'utf8'),
|
||||
ts.ScriptTarget.Latest,
|
||||
true,
|
||||
ts.ScriptKind.TSX,
|
||||
);
|
||||
const code = source.getFullText();
|
||||
|
||||
test('网关列表展示最近一次健康检查结果', () => {
|
||||
assert.match(code, /健康状态/, '应有健康状态列');
|
||||
assert.match(code, /dataIndex: 'health_status'/, '应绑定 health_status 字段');
|
||||
assert.match(code, /gateway\.last_check_reason/, '应展示最近一次检查原因');
|
||||
assert.match(code, /gateway\.last_checked_at/, '应展示最近一次检查时间');
|
||||
assert.match(code, /online|在线/, '健康应展示为在线');
|
||||
assert.match(code, /offline|离线/, '不健康应展示为离线');
|
||||
});
|
||||
|
||||
test('健康状态未检查时不误报异常', () => {
|
||||
// 未检查(unchecked)必须落在独立的静默分支,不允许与离线混用同一个颜色/文案。
|
||||
const render = code.match(/function HealthCell[\s\S]*?\n}/);
|
||||
assert.ok(render, '应有独立的 HealthCell 组件');
|
||||
assert.match(render[0], /未检查/, '应存在未检查文案');
|
||||
const branches = render[0].match(/status === '[a-z]+'/g) ?? [];
|
||||
assert.ok(branches.includes("status === 'healthy'"), '应判断 healthy');
|
||||
assert.ok(branches.includes("status === 'unhealthy'"), '应判断 unhealthy');
|
||||
});
|
||||
|
||||
test('检查原因与时间收进 Tooltip,不挤占表格密度', () => {
|
||||
const render = code.match(/function HealthCell[\s\S]*?\n}/);
|
||||
assert.ok(render, '应有独立的 HealthCell 组件');
|
||||
assert.match(render[0], /Tooltip/, '原因与时间应放 Tooltip');
|
||||
assert.doesNotMatch(render[0], /Typography\.Text type="secondary"/, '不应在单元格内平铺长文本');
|
||||
});
|
||||
Reference in New Issue
Block a user