From e12e03b3b4e6fc102a86941cab2c4aad7b889b93 Mon Sep 17 00:00:00 2001 From: Rogee Date: Wed, 30 Sep 2026 14:33:31 +0800 Subject: [PATCH] =?UTF-8?q?feat(gateway):=20=E5=88=9B=E5=BB=BA/=E7=BC=96?= =?UTF-8?q?=E8=BE=91=E7=BD=91=E5=85=B3=E5=90=8E=E7=AB=8B=E5=8D=B3=E5=BC=82?= =?UTF-8?q?=E6=AD=A5=E6=8E=A2=E6=B4=BB=E5=81=A5=E5=BA=B7=E7=8A=B6=E6=80=81?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - POST/PUT /api/gateways 成功后后台立即探活一次 /healthz 并覆盖落库,不阻塞请求 - probeAndRecordGateway 抽出单网关探活+落库复用;探活超时改为 var 供测试注入 - 网关页保存成功后延迟 2.5s 刷新一次,尽快展示最新健康状态 - 测试:立即探活、慢网关不阻塞、不可达落 unhealthy --- internal/controlplane/api/environments.go | 3 + internal/controlplane/api/gateway_health.go | 35 ++-- .../controlplane/api/gateway_health_test.go | 151 ++++++++++++++++++ internal/controlplane/api/hub_test.go | 31 +++- web/src/pages/gateways/index.tsx | 12 +- 5 files changed, 217 insertions(+), 15 deletions(-) diff --git a/internal/controlplane/api/environments.go b/internal/controlplane/api/environments.go index eb6b7cb..e1e32f9 100644 --- a/internal/controlplane/api/environments.go +++ b/internal/controlplane/api/environments.go @@ -26,6 +26,7 @@ type HubStore interface { ListGateways(ctx context.Context) ([]hub.Gateway, error) GetGateway(ctx context.Context, name string) (hub.Gateway, error) DeleteGateway(ctx context.Context, name string) error + RecordGatewayCheck(ctx context.Context, name, status, reason string) (hub.Gateway, error) ListEnvs(ctx context.Context) ([]hub.Env, error) GetEnv(ctx context.Context, alias string) (hub.Env, error) CreateNetworkExit(ctx context.Context, exit hub.NetworkExit) (hub.NetworkExit, error) @@ -488,6 +489,7 @@ func registerHubWithNetwork(app *fiber.App, store HubStore, probe NetworkExitPro if err != nil { return hubError(c, err) } + probeGatewayAsync(c.Context(), store, gateway) return c.Status(fiber.StatusCreated).JSON(gateway) }) app.Put("/api/gateways/:name", func(c fiber.Ctx) error { @@ -503,6 +505,7 @@ func registerHubWithNetwork(app *fiber.App, store HubStore, probe NetworkExitPro if err != nil { return hubError(c, err) } + probeGatewayAsync(c.Context(), store, gateway) return c.JSON(gateway) }) app.Delete("/api/gateways/:name", func(c fiber.Ctx) error { diff --git a/internal/controlplane/api/gateway_health.go b/internal/controlplane/api/gateway_health.go index 296986d..6c87300 100644 --- a/internal/controlplane/api/gateway_health.go +++ b/internal/controlplane/api/gateway_health.go @@ -10,7 +10,8 @@ import ( "github.com/sirupsen/logrus" ) -const gatewayHealthProbeTimeout = 5 * time.Second +// gatewayHealthProbeTimeout 是单次 /healthz 探活的超时;测试可临时替换。 +var gatewayHealthProbeTimeout = 5 * time.Second // gatewayHealthStore 是网关健康检查所需的存储能力;生产实现为 *hub.Store,测试使用内存桩。 type gatewayHealthStore interface { @@ -44,7 +45,6 @@ func RecordGatewayHealthChecks(ctx context.Context, store *hub.Store) error { } // recordGatewayHealthChecks 对全部注册网关并发探活一次,每个网关只覆盖写入最近一次结果。 -// 单个网关落库失败不阻断其余网关,仅记日志(下一轮探活会重试)。 func recordGatewayHealthChecks(ctx context.Context, store gatewayHealthStore) error { gateways, err := store.ListGateways(ctx) if err != nil { @@ -55,18 +55,29 @@ func recordGatewayHealthChecks(ctx context.Context, store gatewayHealthStore) er 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") - } + probeAndRecordGateway(ctx, store, gateway) }() } wg.Wait() return nil } + +// probeAndRecordGateway 对单个网关探活一次并覆盖写入最近一次结果;落库失败不阻断其余网关,仅记日志。 +func probeAndRecordGateway(ctx context.Context, store gatewayHealthStore, gateway hub.Gateway) { + 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") + } +} + +// probeGatewayAsync 在后台对单个网关立即探活一次并落库,不阻塞当前请求; +// 用 WithoutCancel 脱离请求生命周期,请求结束后探活与落库仍会完成。 +func probeGatewayAsync(ctx context.Context, store gatewayHealthStore, gateway hub.Gateway) { + go probeAndRecordGateway(context.WithoutCancel(ctx), store, gateway) +} diff --git a/internal/controlplane/api/gateway_health_test.go b/internal/controlplane/api/gateway_health_test.go index 099aa5f..2fcccc9 100644 --- a/internal/controlplane/api/gateway_health_test.go +++ b/internal/controlplane/api/gateway_health_test.go @@ -10,8 +10,21 @@ import ( "time" hub "git.ipao.vip/rogee/creator-hub/internal/environment" + "github.com/gofiber/fiber/v3" ) +// setGatewayHealthProbeTimeout 临时替换探活超时,返回旧值供 restore 恢复。 +func setGatewayHealthProbeTimeout(d time.Duration) time.Duration { + old := gatewayHealthProbeTimeout + gatewayHealthProbeTimeout = d + return old +} + +// restoreGatewayHealthProbeTimeout 恢复探活超时旧值。 +func restoreGatewayHealthProbeTimeout(old time.Duration) func() { + return func() { gatewayHealthProbeTimeout = old } +} + type recordedGatewayCheck struct { name string status string @@ -120,3 +133,141 @@ func TestGatewayHealthProbeTimeoutBounded(t *testing.T) { t.Fatalf("probe timeout %v must stay below scheduling cadence", gatewayHealthProbeTimeout) } } + +// waitForGatewayCheck 轮询等待指定网关的立即探活结果落库;超时判失败。 +func waitForGatewayCheck(t *testing.T, store *healthMemoryStore, name string) recordedGatewayCheck { + t.Helper() + deadline := time.Now().Add(2 * time.Second) + for time.Now().Before(deadline) { + store.mu.Lock() + for _, check := range store.checks { + if check.name == name { + store.mu.Unlock() + return check + } + } + store.mu.Unlock() + time.Sleep(5 * time.Millisecond) + } + t.Fatalf("gateway %q health check was not recorded within deadline", name) + return recordedGatewayCheck{} +} + +// healthzProbeServer 返回校验 Bearer token 后立即响应 204 的 /healthz 测试服务。 +func healthzProbeServer(t *testing.T, token string, probeStarted chan struct{}, respond func(http.ResponseWriter)) *httptest.Server { + t.Helper() + server := httptest.NewServer(http.HandlerFunc(func(response http.ResponseWriter, request *http.Request) { + if request.URL.Path != "/healthz" { + response.WriteHeader(http.StatusNotFound) + return + } + if request.Header.Get("Authorization") != "Bearer "+token { + response.WriteHeader(http.StatusUnauthorized) + return + } + if probeStarted != nil { + close(probeStarted) + } + if respond != nil { + respond(response) + return + } + response.WriteHeader(http.StatusNoContent) + })) + t.Cleanup(server.Close) + return server +} + +// TestCreateGatewayProbesHealthImmediately:注册网关成功后立即按注册的 endpoint 与 token +// 对 /healthz 探活并落库,结果为 healthy。 +func TestCreateGatewayProbesHealthImmediately(t *testing.T) { + store := &healthMemoryStore{memoryStore: newMemoryStore()} + probeStarted := make(chan struct{}) + server := healthzProbeServer(t, "unit-test-gateway-token", probeStarted, nil) + app := fiber.New() + registerHubWithNetwork(app, store, fakeExitProbe{}, func(hub.NetworkExitAccess) (string, error) { return "", nil }) + + response := do(app, http.MethodPost, "/api/gateways", + `{"name":"gw-new","endpoint":"`+server.URL+`","token":"unit-test-gateway-token"}`) + if response.Code != http.StatusCreated { + t.Fatalf("create gateway returned %d: %s", response.Code, response.Body.String()) + } + + select { + case <-probeStarted: + case <-time.After(2 * time.Second): + t.Fatal("gateway health probe was not triggered after create") + } + check := waitForGatewayCheck(t, store, "gw-new") + if check.status != "healthy" || check.reason != "" { + t.Fatalf("created gateway check = %#v, want healthy without reason", check) + } +} + +// TestCreateGatewayDoesNotBlockOnSlowGateway:网关探活被拖住时注册请求必须立即返回, +// 后台探活超时后把 unhealthy 结果落库。 +func TestCreateGatewayDoesNotBlockOnSlowGateway(t *testing.T) { + t.Cleanup(restoreGatewayHealthProbeTimeout(setGatewayHealthProbeTimeout(250 * time.Millisecond))) + store := &healthMemoryStore{memoryStore: newMemoryStore()} + server := healthzProbeServer(t, "unit-test-gateway-token", nil, func(response http.ResponseWriter) { + time.Sleep(800 * time.Millisecond) // 超过探活超时,若探活同步执行则注册请求会被拖住 + response.WriteHeader(http.StatusNoContent) + }) + app := fiber.New() + registerHubWithNetwork(app, store, fakeExitProbe{}, func(hub.NetworkExitAccess) (string, error) { return "", nil }) + + started := time.Now() + response := do(app, http.MethodPost, "/api/gateways", + `{"name":"gw-slow","endpoint":"`+server.URL+`","token":"unit-test-gateway-token"}`) + if response.Code != http.StatusCreated { + t.Fatalf("create gateway returned %d: %s", response.Code, response.Body.String()) + } + if elapsed := time.Since(started); elapsed > 500*time.Millisecond { + t.Fatalf("create gateway blocked %v on health probe", elapsed) + } + check := waitForGatewayCheck(t, store, "gw-slow") + if check.status != "unhealthy" || check.reason == "" { + t.Fatalf("slow gateway check = %#v, want unhealthy with reason", check) + } +} + +// TestUpdateGatewayProbesHealthImmediately:编辑网关成功后立即按新 endpoint 探活并覆盖落库。 +func TestUpdateGatewayProbesHealthImmediately(t *testing.T) { + store := &healthMemoryStore{memoryStore: newMemoryStore()} + store.gateways["gw-1"] = hub.Gateway{Name: "gw-1", Endpoint: "http://gw-old:8081", Token: "unit-test-gateway-token"} + server := healthzProbeServer(t, "unit-test-gateway-token", nil, nil) + app := fiber.New() + registerHubWithNetwork(app, store, fakeExitProbe{}, func(hub.NetworkExitAccess) (string, error) { return "", nil }) + + response := do(app, http.MethodPut, "/api/gateways/gw-1", + `{"name":"gw-1","endpoint":"`+server.URL+`","token":""}`) + if response.Code != http.StatusOK { + t.Fatalf("update gateway returned %d: %s", response.Code, response.Body.String()) + } + check := waitForGatewayCheck(t, store, "gw-1") + if check.status != "healthy" || check.reason != "" { + t.Fatalf("updated gateway check = %#v, want healthy without reason", check) + } +} + +// TestUpdateGatewayUnreachableRecordsUnhealthy:编辑指向不可达地址时保存照常成功, +// 后台探活把 unhealthy 结果落库供页面展示。 +func TestUpdateGatewayUnreachableRecordsUnhealthy(t *testing.T) { + t.Cleanup(restoreGatewayHealthProbeTimeout(setGatewayHealthProbeTimeout(300*time.Millisecond))) + store := &healthMemoryStore{memoryStore: newMemoryStore()} + store.gateways["gw-1"] = hub.Gateway{Name: "gw-1", Endpoint: "http://gw-old:8081", Token: "unit-test-gateway-token"} + closedServer := httptest.NewServer(http.HandlerFunc(func(http.ResponseWriter, *http.Request) {})) + closedServer.Close() // 端口已释放,保证连接失败 + app := fiber.New() + registerHubWithNetwork(app, store, fakeExitProbe{}, func(hub.NetworkExitAccess) (string, error) { return "", nil }) + + response := do(app, http.MethodPut, "/api/gateways/gw-1", + `{"name":"gw-1","endpoint":"`+closedServer.URL+`","token":""}`) + if response.Code != http.StatusOK { + t.Fatalf("update gateway returned %d: %s", response.Code, response.Body.String()) + } + check := waitForGatewayCheck(t, store, "gw-1") + if check.status != "unhealthy" || check.reason == "" { + t.Fatalf("unreachable gateway check = %#v, want unhealthy with reason", check) + } +} diff --git a/internal/controlplane/api/hub_test.go b/internal/controlplane/api/hub_test.go index cec27e6..4c5e152 100644 --- a/internal/controlplane/api/hub_test.go +++ b/internal/controlplane/api/hub_test.go @@ -132,8 +132,35 @@ func (s *memoryStore) lock(key string) func() { return lock.Unlock } -func (s *memoryStore) CreateGateway(_ context.Context, _, _, _ string) (hub.Gateway, error) { - return hub.Gateway{}, nil +func (s *memoryStore) CreateGateway(_ context.Context, name, endpoint, token string) (hub.Gateway, error) { + s.mu.Lock() + defer s.mu.Unlock() + if name == "" || token == "" { + return hub.Gateway{}, hub.ErrInvalid + } + if s.gateways == nil { + s.gateways = map[string]hub.Gateway{} + } + if _, exists := s.gateways[name]; exists { + return hub.Gateway{}, hub.ErrConflict + } + gateway := hub.Gateway{Name: name, Endpoint: endpoint, Token: token} + s.gateways[name] = gateway + return gateway, nil +} + +func (s *memoryStore) RecordGatewayCheck(_ context.Context, name, status, reason string) (hub.Gateway, error) { + s.mu.Lock() + defer s.mu.Unlock() + gateway, exists := s.gateways[name] + if !exists { + return hub.Gateway{}, hub.ErrNotFound + } + gateway.HealthStatus, gateway.LastCheckReason = status, reason + now := time.Now() + gateway.LastCheckedAt = &now + s.gateways[name] = gateway + return gateway, nil } func (s *memoryStore) UpdateGateway(_ context.Context, currentName, name, endpoint, token string) (hub.Gateway, error) { s.mu.Lock() diff --git a/web/src/pages/gateways/index.tsx b/web/src/pages/gateways/index.tsx index d11282d..0237347 100644 --- a/web/src/pages/gateways/index.tsx +++ b/web/src/pages/gateways/index.tsx @@ -1,5 +1,5 @@ // 网关管理:语义对齐 web.archived GatewaysPage.jsx(注册/编辑 Modal、令牌显隐复制、删除)。 -import { useCallback, useEffect, useState } from 'react'; +import { useCallback, useEffect, useRef, useState } from 'react'; 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'; @@ -96,6 +96,7 @@ export default function Page() { const [created, setCreated] = useState(null); const { message: messageApi } = App.useApp(); const { vertical, horizontal } = useOverflowGrid(gateways.length); + const healthReloadTimer = useRef>(undefined); const load = useCallback(async () => { setPending(true); @@ -110,6 +111,13 @@ export default function Page() { } }, []); + // 创建/编辑后的立即探活在后台异步进行:延迟刷新一次让健康状态列尽快显示最新结果。 + const scheduleHealthReload = useCallback(() => { + clearTimeout(healthReloadTimer.current); + healthReloadTimer.current = setTimeout(load, 2500); + }, [load]); + useEffect(() => () => clearTimeout(healthReloadTimer.current), []); + useEffect(() => { load(); }, [load]); @@ -120,6 +128,7 @@ export default function Page() { const record = await create('gateways', values); setCreated(record?.data ?? record); await load(); + scheduleHealthReload(); setCreateOpen(false); return true; } catch (reason) { @@ -136,6 +145,7 @@ export default function Page() { try { await update('gateways', editing.name, values); await load(); + scheduleHealthReload(); setEditing(null); return true; } catch (reason) {