feat(ha): 网关多副本安全 —— 事件订阅改队列组,杜绝重复落库/重复计费
去单点的代码层基础:dispatcher/mcp-go 本就靠队列组可多副本,但网关侧的 eval/usage/ status/config 订阅是广播(nc.Subscribe),多网关副本下每条会被每个副本各处理一遍 → 评测/状态重复写 PG、token 用量重复累加(日预算翻倍)、config 请求多份重复应答。 改为 QueueSubscribe + contract.QueueGateway 队列组:组内每条事件/请求只一个副本处理。 (dispatcher/mcp-go 的 config 变更广播订阅保持 nc.Subscribe 不动——每副本都要热更新。) 验证: - 单测 TestGatewayQueueDedup:2 网关副本订阅,发 50 条 eval,合计处理 50 次(非 100)。 - live:2 dispatcher 副本提 8 任务,队列组自动 4/4 分摊。 至此进程级全部可水平复制。剩 NATS 集群 / 网关 LB / PG·Redis·Milvus 基础设施 HA 属部署期。 Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
This commit is contained in:
@@ -333,8 +333,9 @@ func (b *Bus) PublishTaskStatus(ev *contract.TaskStatusEvent) error {
|
||||
}
|
||||
|
||||
// SubscribeTaskStatus 订阅任务状态流转(网关调用,落 PG + 推 UI)。
|
||||
// 队列组订阅:多网关副本下每条状态只由一个副本落库,避免重复写(HA)。
|
||||
func (b *Bus) SubscribeTaskStatus(onEvent func(*contract.TaskStatusEvent)) (unsub func() error, err error) {
|
||||
sub, err := b.nc.Subscribe(contract.SubjectTaskStatus, func(m *nats.Msg) {
|
||||
sub, err := b.nc.QueueSubscribe(contract.SubjectTaskStatus, contract.QueueGateway, func(m *nats.Msg) {
|
||||
var ev contract.TaskStatusEvent
|
||||
if json.Unmarshal(m.Data, &ev) == nil {
|
||||
onEvent(&ev)
|
||||
@@ -357,9 +358,9 @@ func (b *Bus) PublishEval(ev *contract.EvalEvent) error {
|
||||
return b.nc.Publish(contract.SubjectEval, data)
|
||||
}
|
||||
|
||||
// SubscribeEval 订阅评测结果(网关调用,落 PG)。
|
||||
// SubscribeEval 订阅评测结果(网关调用,落 PG)。队列组:多副本下每条只落一次(HA)。
|
||||
func (b *Bus) SubscribeEval(onEvent func(*contract.EvalEvent)) (unsub func() error, err error) {
|
||||
sub, err := b.nc.Subscribe(contract.SubjectEval, func(m *nats.Msg) {
|
||||
sub, err := b.nc.QueueSubscribe(contract.SubjectEval, contract.QueueGateway, func(m *nats.Msg) {
|
||||
var ev contract.EvalEvent
|
||||
if json.Unmarshal(m.Data, &ev) == nil {
|
||||
onEvent(&ev)
|
||||
@@ -381,8 +382,9 @@ func (b *Bus) PublishUsage(ev *contract.UsageEvent) error {
|
||||
}
|
||||
|
||||
// SubscribeUsage 订阅 token 用量(网关调用,累计到用户日预算 / 计费)。
|
||||
// 队列组:多副本下每条用量只累加一次,避免日预算被重复计(HA)。
|
||||
func (b *Bus) SubscribeUsage(onEvent func(*contract.UsageEvent)) (unsub func() error, err error) {
|
||||
sub, err := b.nc.Subscribe(contract.SubjectUsage, func(m *nats.Msg) {
|
||||
sub, err := b.nc.QueueSubscribe(contract.SubjectUsage, contract.QueueGateway, func(m *nats.Msg) {
|
||||
var ev contract.UsageEvent
|
||||
if json.Unmarshal(m.Data, &ev) == nil {
|
||||
onEvent(&ev)
|
||||
@@ -480,8 +482,9 @@ func (b *Bus) RequestConfigWithRetry(ctx context.Context, kind string, apply fun
|
||||
}
|
||||
|
||||
// ServeConfig 让控制面响应某 kind 的配置请求;provide 返回当前激活配置(可为 nil)。
|
||||
// 队列组:多网关副本下每个配置请求只由一个副本应答,避免请求方收到多份重复应答(HA)。
|
||||
func (b *Bus) ServeConfig(kind string, provide func() *contract.ModelConfig) (unsub func() error, err error) {
|
||||
sub, err := b.nc.Subscribe(contract.ConfigGetSubject(kind), func(m *nats.Msg) {
|
||||
sub, err := b.nc.QueueSubscribe(contract.ConfigGetSubject(kind), contract.QueueGateway, func(m *nats.Msg) {
|
||||
var data []byte
|
||||
if cfg := provide(); cfg != nil {
|
||||
data, _ = json.Marshal(cfg)
|
||||
|
||||
@@ -3,6 +3,7 @@ package bus_test
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"sync/atomic"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
@@ -360,3 +361,48 @@ func TestGracefulDrain(t *testing.T) {
|
||||
}
|
||||
t.Logf("✓ drain 等待在途任务跑完:耗时 %v", elapsed)
|
||||
}
|
||||
|
||||
// TestGatewayQueueDedup 验证多网关副本下事件不重复处理:两个网关实例都订阅评测/用量/状态,
|
||||
// 经队列组应「每条只被一个副本处理」(旧的广播订阅会被两副本各处理一遍 → 重复落库/重复计费)。
|
||||
func TestGatewayQueueDedup(t *testing.T) {
|
||||
url := startEmbeddedNATS(t)
|
||||
pub, err := bus.Connect(url)
|
||||
if err != nil {
|
||||
t.Fatalf("pub connect: %v", err)
|
||||
}
|
||||
defer pub.Close()
|
||||
gwA, err := bus.Connect(url)
|
||||
if err != nil {
|
||||
t.Fatalf("gwA connect: %v", err)
|
||||
}
|
||||
defer gwA.Close()
|
||||
gwB, err := bus.Connect(url)
|
||||
if err != nil {
|
||||
t.Fatalf("gwB connect: %v", err)
|
||||
}
|
||||
defer gwB.Close()
|
||||
|
||||
var total int64
|
||||
count := func(_ *contract.EvalEvent) { atomic.AddInt64(&total, 1) }
|
||||
if _, err := gwA.SubscribeEval(count); err != nil {
|
||||
t.Fatalf("gwA sub: %v", err)
|
||||
}
|
||||
if _, err := gwB.SubscribeEval(count); err != nil {
|
||||
t.Fatalf("gwB sub: %v", err)
|
||||
}
|
||||
time.Sleep(100 * time.Millisecond) // 等订阅就绪
|
||||
|
||||
const n = 50
|
||||
for i := 0; i < n; i++ {
|
||||
if err := pub.PublishEval(&contract.EvalEvent{TaskID: "t", Overall: 1}); err != nil {
|
||||
t.Fatalf("publish: %v", err)
|
||||
}
|
||||
}
|
||||
time.Sleep(400 * time.Millisecond) // 等投递
|
||||
|
||||
got := atomic.LoadInt64(&total)
|
||||
if got != n {
|
||||
t.Fatalf("两副本共处理 %d 条,期望 %d(广播会得 %d=重复处理)", got, n, 2*n)
|
||||
}
|
||||
t.Logf("✓ 队列组去重:%d 条事件被两网关副本合计处理 %d 次(无重复)", n, got)
|
||||
}
|
||||
|
||||
@@ -25,6 +25,10 @@ const (
|
||||
QueueToolsGo = "mcp-go-workers" // mcp-go 队列组(多副本负载均衡)
|
||||
QueueToolsPy = "mcp-py-workers" // mcp-py 队列组
|
||||
|
||||
// QueueGateway 是网关侧事件订阅/配置应答的队列组:多网关副本下,每条
|
||||
// eval/usage/status 事件与每个 config 请求只由组内一个副本处理,避免重复落库/重复应答(HA)。
|
||||
QueueGateway = "gateway-workers"
|
||||
|
||||
// 服务探活:dispatcher 既无 HTTP 端点也不挂工具,单独用一个 core NATS
|
||||
// request-reply 心跳主题让控制面(管理端「服务状态」)能判定它在不在线。
|
||||
SubjectHealthDispatcher = "sundynix.health.dispatcher"
|
||||
|
||||
Reference in New Issue
Block a user