From 05c25d7099bf0c05dc364b0546d8f610b7326a3e Mon Sep 17 00:00:00 2001 From: Blizzard Date: Fri, 26 Jun 2026 10:29:23 +0800 Subject: [PATCH] =?UTF-8?q?feat(ops):=20=E4=BC=98=E9=9B=85=E5=81=9C?= =?UTF-8?q?=E6=9C=BA=20drain=20=E2=80=94=E2=80=94=20=E4=B8=89=20Go=20?= =?UTF-8?q?=E6=9C=8D=E5=8A=A1=20SIGTERM=20=E5=90=8E=E6=8E=92=E7=A9=BA?= =?UTF-8?q?=E5=9C=A8=E9=80=94=EF=BC=8C=E4=B8=8D=E7=A1=AC=E5=88=87?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 滚动更新/重启时旧的"信号到→进程退"会硬切在途工作:dispatcher 在途任务被掐断、 mcp-go 在途工具调用让 dispatcher 干等超时、gateway 在途 HTTP 请求被截断。 共享 bus 加在途追踪 + drain: - ConsumeTasks/ServeTool 各加 sync.WaitGroup 跟踪在途 goroutine,返回 drain(ctx): 先停止接新活(cc.Stop / Unsubscribe),再等在途跑完至 drain 超时。 - 关键修复:任务 handler 的 ctx 改为派生自 context.Background()(而非信号 ctx), 否则 SIGTERM 会立即取消在途任务的 ctx,drain 形同虚设。超时未跑完才由 JetStream AckWait 重投兜底(不丢任务)。 - DrainTimeout():SHUTDOWN_DRAIN_TIMEOUT 秒,默认 30s。 各服务收尾: - gateway:r.Run → http.Server + signal.NotifyContext + srv.Shutdown(drain 在途请求), 随后 defer 关 db/redis/bus(HTTP 排空后才断后端连接)。 - dispatcher:收到信号 → drain 在途任务跑完再退。 - mcp-go:收到信号 → drain 在途工具调用回完再退(dispatcher 拿到结果而非超时)。 bus 加 TestGracefulDrain(drain 等满在途任务 + 验证在途 ctx 不被取消),e2e 测试适配 新签名,四模块全绿。live:在途任务执行中 kill -TERM dispatcher → 日志「drain 在途任务」 → 11s 后 task done(511字完整生成) → 「drain 完成,退出」,任务终态 done 非 failed; gateway/mcp-go 同样优雅退出。 Co-Authored-By: Claude Opus 4.8 (1M context) --- project_analysis.md | 4 +- .../internal/nats/subscriber.go | 14 ++-- sundynix-gateway/cmd/server/main.go | 29 +++++++- sundynix-mcp-go/internal/mcp/gateway.go | 14 ++-- sundynix-shared/bus/bus.go | 49 +++++++++++-- sundynix-shared/bus/bus_e2e_test.go | 70 ++++++++++++++++--- 6 files changed, 153 insertions(+), 27 deletions(-) diff --git a/project_analysis.md b/project_analysis.md index d45231f..8977289 100644 --- a/project_analysis.md +++ b/project_analysis.md @@ -147,7 +147,7 @@ Harness = 围绕 LLM 的可靠性 / 安全 / 质量治理层。4 个组件均为 | **备份 / 灾备** | PG/Milvus/Neo4j 无备份恢复方案 | 数据丢失风险 | | **安全审计** | 全自述,未渗透/未审计 | "纵深防御"未被验证 | | **配额 / 多租户** | 仅按 IP 全局限流;owner 基础隔离 | 无团队/组织/按用户配额 | -| **可靠性细节** | exec 轨迹流丢事件;停机不 drain | 观测缺口 + 滚更硬切在途任务 | +| **可靠性细节** | exec 轨迹流丢事件;~~停机不 drain~~(已优雅停机 drain) | 观测缺口 | | **计量计费** | 计价配置有,计量×单价+配额未落地 | 商业化未闭环 | | **本地模型** | 仅 OpenAI 兼容在线 API | vLLM/Ollama、模型路由/fallback 未做 | @@ -207,7 +207,7 @@ Harness = 围绕 LLM 的可靠性 / 安全 / 质量治理层。4 个组件均为 | P0 | 本地模型(vLLM/Ollama) + 推理模型 reasoning_content 适配 | 对齐生产 Qwen | | P1 | 高可用:网关/调度多副本 + NATS 集群 + 自愈 | 解单点 | | P1 | 备份/灾备演练(PG/Milvus/Neo4j) | 数据安全 | -| P1 | exec 轨迹 Redis 回放 + 优雅停机 drain | 可靠性 | +| ~~P1~~ ✅ | ~~优雅停机 drain~~:三 Go 服务全覆盖(gateway HTTP Shutdown / dispatcher 在途任务跑完 / mcp-go 在途工具回完),SIGTERM 后等在途至 SHUTDOWN_DRAIN_TIMEOUT(默认30s)再退;在途任务 ctx 脱离信号 ctx 不被掐断。剩 exec 轨迹 Redis 回放 | 可靠性 | | P1 | 核心链路集成测试 + 前端 E2E | 测试纵深 | | P2 | 计量计费闭环 + 按用户配额;OCR 完善;KbView 拆分;循环节点 | | | P2 | 安全审计 / 渗透(让"纵深防御"从自述变已验证) | 外部 | diff --git a/sundynix-dispatcher/internal/nats/subscriber.go b/sundynix-dispatcher/internal/nats/subscriber.go index 8a2a1c0..6bf0c5d 100644 --- a/sundynix-dispatcher/internal/nats/subscriber.go +++ b/sundynix-dispatcher/internal/nats/subscriber.go @@ -33,15 +33,21 @@ func MustConnect(url string) *Subscriber { // ConsumeTasks 从 sundynix.tasks.* 持续消费任务(队列组负载均衡),阻塞至 ctx 取消。 func (s *Subscriber) ConsumeTasks(ctx context.Context, h TaskHandler) error { - stop, err := s.inner.ConsumeTasks(ctx, func(c context.Context, t *contract.Task) error { + drain, err := s.inner.ConsumeTasks(ctx, func(c context.Context, t *contract.Task) error { return h(c, t) }) if err != nil { return err } - defer stop() - <-ctx.Done() - return ctx.Err() + <-ctx.Done() // 收到停机信号 + // 优雅停机:停止消费新任务 + 等在途任务跑完(至多 DrainTimeout),再退出。 + to := sharedbus.DrainTimeout() + log.Printf("[dispatcher] 收到停机信号,drain 在途任务(≤%s)…", to) + dctx, cancel := context.WithTimeout(context.Background(), to) + defer cancel() + drain(dctx) + log.Printf("[dispatcher] drain 完成,退出") + return nil } // PublishToken / CompleteStream 让 Subscriber 满足 eino.TokenSink, diff --git a/sundynix-gateway/cmd/server/main.go b/sundynix-gateway/cmd/server/main.go index 9eda2e3..53bf00f 100644 --- a/sundynix-gateway/cmd/server/main.go +++ b/sundynix-gateway/cmd/server/main.go @@ -5,13 +5,17 @@ import ( "context" "encoding/json" "log" + "net/http" "os" + "os/signal" + "syscall" "time" "github.com/sundynix/sundynix-gateway/internal/blob" "github.com/sundynix/sundynix-gateway/internal/nats" "github.com/sundynix/sundynix-gateway/internal/router" "github.com/sundynix/sundynix-gateway/internal/store" + sharedbus "github.com/sundynix/sundynix-shared/bus" "github.com/sundynix/sundynix-shared/contract" "github.com/sundynix/sundynix-shared/otelx" "github.com/sundynix/sundynix-shared/secrets" @@ -94,10 +98,29 @@ func main() { r := router.New(db, cache, bus, blobStore) addr := envOr("GATEWAY_ADDR", ":8080") - log.Printf("[gateway] listening on %s", addr) - if err := r.Run(addr); err != nil { - log.Fatalf("[gateway] exit: %v", err) + srv := &http.Server{Addr: addr, Handler: r} + + // 后台监听;ListenAndServe 在 Shutdown 后返回 ErrServerClosed(正常退出)。 + go func() { + log.Printf("[gateway] listening on %s", addr) + if err := srv.ListenAndServe(); err != nil && err != http.ErrServerClosed { + log.Fatalf("[gateway] listen: %v", err) + } + }() + + // 优雅停机:收到信号 → 停止收新连接、drain 在途请求(至多 DrainTimeout),再退出 + // (defer 的 db/cache/bus/blob Close 随后执行,确保 HTTP 排空后才断后端连接)。 + ctx, stop := signal.NotifyContext(context.Background(), syscall.SIGINT, syscall.SIGTERM) + defer stop() + <-ctx.Done() + to := sharedbus.DrainTimeout() + log.Printf("[gateway] 收到停机信号,优雅关闭 HTTP(drain 在途请求 ≤%s)…", to) + sctx, cancel := context.WithTimeout(context.Background(), to) + defer cancel() + if err := srv.Shutdown(sctx); err != nil { + log.Printf("[gateway] HTTP 优雅关闭超时/出错: %v", err) } + log.Println("[gateway] 已优雅停机") } func envOr(key, def string) string { diff --git a/sundynix-mcp-go/internal/mcp/gateway.go b/sundynix-mcp-go/internal/mcp/gateway.go index 78ea6bb..6d64931 100644 --- a/sundynix-mcp-go/internal/mcp/gateway.go +++ b/sundynix-mcp-go/internal/mcp/gateway.go @@ -68,15 +68,21 @@ func NewGateway(b *sharedbus.Bus, s *search.Hybrid, m *memory.Store, h *history. // Serve 以队列组通配订阅 sundynix.tools.go.>,按工具名分发并阻塞。 func (g *Gateway) Serve(ctx context.Context) error { - unsub, err := g.bus.ServeTool(contract.SubjectToolsGoAll, contract.QueueToolsGo, g.dispatch) + drain, err := g.bus.ServeTool(contract.SubjectToolsGoAll, contract.QueueToolsGo, g.dispatch) if err != nil { return err } - defer func() { _ = unsub() }() log.Printf("[mcp_go] tools ready on %s (queue=%s): wiki_search, kb_ingest, kb_search, kb_graph, report_render, memory_*, history_*, echo", contract.SubjectToolsGoAll, contract.QueueToolsGo) - <-ctx.Done() - return ctx.Err() + <-ctx.Done() // 收到停机信号 + // 优雅停机:停止接新工具调用 + 等在途调用回完(至多 DrainTimeout),让 dispatcher 拿到结果而非干等超时。 + to := sharedbus.DrainTimeout() + log.Printf("[mcp_go] 收到停机信号,drain 在途工具调用(≤%s)…", to) + dctx, cancel := context.WithTimeout(context.Background(), to) + defer cancel() + drain(dctx) + log.Printf("[mcp_go] drain 完成,退出") + return nil } // buildRegistry 注册 mcp-go 全部工具:名称 → (中文名, 作用, 处理函数)。 diff --git a/sundynix-shared/bus/bus.go b/sundynix-shared/bus/bus.go index a8c3052..1a03223 100644 --- a/sundynix-shared/bus/bus.go +++ b/sundynix-shared/bus/bus.go @@ -9,6 +9,7 @@ import ( "log" "os" "strconv" + "sync" "time" "github.com/nats-io/nats.go" @@ -246,14 +247,17 @@ func toolConcurrency() int { // ServeTool 以队列组订阅工具主题(可用通配 sundynix.tools.go.>),队列组内多副本自动负载均衡(水平扩)。 // 单实例内每个请求分发到独立 goroutine 并发处理(垂直扩,上限 MCP_TOOL_CONCURRENCY)—— // 回调立即返回不阻塞 NATS 投递;信号量约束同时执行的 handler 数。返回的 unsub 用于退订。 -func (b *Bus) ServeTool(subject, queue string, h ToolHandler) (unsub func() error, err error) { +func (b *Bus) ServeTool(subject, queue string, h ToolHandler) (drain func(context.Context), err error) { sem := make(chan struct{}, toolConcurrency()) + var wg sync.WaitGroup // 跟踪在途工具调用,供优雅停机 drain 等待回完 sub, err := b.nc.QueueSubscribe(subject, queue, func(m *nats.Msg) { // 立即起 goroutine 并返回,让 NATS 继续投递下一条(核心 NATS 单订阅回调本是串行)。 + wg.Add(1) go func() { sem <- struct{}{} // 限并发:超额则在此短暂排队 defer func() { <-sem + wg.Done() if r := recover(); r != nil { // handler panic → 回错误结果,调用方不必干等超时 respond(m, &contract.ToolResult{OK: false, Error: fmt.Sprintf("tool panic: %v", r)}) } @@ -279,7 +283,11 @@ func (b *Bus) ServeTool(subject, queue string, h ToolHandler) (unsub func() erro if err != nil { return nil, fmt.Errorf("serve tool %s: %w", subject, err) } - return sub.Unsubscribe, nil + // drain:退订(停止接新请求),等在途工具调用回完(至多 dctx)。 + return func(dctx context.Context) { + _ = sub.Unsubscribe() + drainWait(&wg, dctx) + }, nil } func respond(m *nats.Msg, res *contract.ToolResult) { @@ -523,10 +531,31 @@ func taskConcurrency() int { return 8 } +// DrainTimeout 是优雅停机时等待在途工作(任务/工具调用/HTTP 请求)跑完的上限。 +// 经 SHUTDOWN_DRAIN_TIMEOUT 秒配置,缺省 30s。超时即放弃等待退出(在途未 ack 由 JetStream 重投兜底)。 +func DrainTimeout() time.Duration { + if v := os.Getenv("SHUTDOWN_DRAIN_TIMEOUT"); v != "" { + if n, err := strconv.Atoi(v); err == nil && n > 0 { + return time.Duration(n) * time.Second + } + } + return 30 * time.Second +} + +// drainWait 停止接新活后,等待在途 WaitGroup 跑完,或 dctx 到期放弃。 +func drainWait(wg *sync.WaitGroup, dctx context.Context) { + done := make(chan struct{}) + go func() { wg.Wait(); close(done) }() + select { + case <-done: + case <-dctx.Done(): + } +} + // ConsumeTasks 在持久消费者上消费任务,队列组内负载均衡。 // 每个任务分发到独立 worker goroutine 并发执行——一个慢任务/HITL 待审不再阻塞后续任务。 // 并发上限由信号量 + 消费者 MaxAckPending 双重约束(背压)。返回的 stop 用于优雅停止消费。 -func (b *Bus) ConsumeTasks(ctx context.Context, h TaskHandler) (stop func(), err error) { +func (b *Bus) ConsumeTasks(ctx context.Context, h TaskHandler) (drain func(context.Context), err error) { concurrency := taskConcurrency() cons, err := b.js.CreateOrUpdateConsumer(ctx, contract.StreamTasks, jetstream.ConsumerConfig{ Durable: contract.ConsumerDurable, @@ -542,6 +571,7 @@ func (b *Bus) ConsumeTasks(ctx context.Context, h TaskHandler) (stop func(), err return nil, fmt.Errorf("create consumer: %w", err) } sem := make(chan struct{}, concurrency) // 限并发:最多 N 个任务同时执行 + var wg sync.WaitGroup // 跟踪在途任务,供优雅停机 drain 等待跑完 cc, err := cons.Consume(func(msg jetstream.Msg) { t, err := contract.Unmarshal(msg.Data()) if err != nil { @@ -554,17 +584,20 @@ func (b *Bus) ConsumeTasks(ctx context.Context, h TaskHandler) (stop func(), err case <-ctx.Done(): return } + wg.Add(1) go func() { defer func() { <-sem // 释放并发额度 + wg.Done() if r := recover(); r != nil { // 任务处理 panic:丢弃不重投(避免崩溃循环),记录后继续。 log.Printf("[bus] task %s handler panic: %v", t.ID, r) _ = msg.Term() } }() - // 从消息头还原上游链路,开消费 span(成为 gateway 发布 span 的子节点)。 - mctx := extractTrace(ctx, nats.Header(msg.Headers())) + // 关键:handler ctx 派生自 Background(而非信号 ctx),使「停止消费」不会立刻掐断在途任务—— + // 在途任务由 drain 等待至完成;drain 超时未跑完才由 JetStream AckWait 重投兜底。 + mctx := extractTrace(context.Background(), nats.Header(msg.Headers())) mctx, span := tracer().Start(mctx, "nats.consume task", trace.WithSpanKind(trace.SpanKindConsumer), trace.WithAttributes(attribute.String("sundynix.task_id", t.ID))) @@ -580,5 +613,9 @@ func (b *Bus) ConsumeTasks(ctx context.Context, h TaskHandler) (stop func(), err if err != nil { return nil, fmt.Errorf("consume: %w", err) } - return cc.Stop, nil + // drain:停止新投递,等在途任务跑完(至多 dctx)。 + return func(dctx context.Context) { + cc.Stop() + drainWait(&wg, dctx) + }, nil } diff --git a/sundynix-shared/bus/bus_e2e_test.go b/sundynix-shared/bus/bus_e2e_test.go index c5f6bde..fbe31ea 100644 --- a/sundynix-shared/bus/bus_e2e_test.go +++ b/sundynix-shared/bus/bus_e2e_test.go @@ -55,14 +55,14 @@ func TestTaskRoundTrip(t *testing.T) { defer dp.Close() got := make(chan *contract.Task, 1) - stop, err := dp.ConsumeTasks(ctx, func(_ context.Context, task *contract.Task) error { + drain, err := dp.ConsumeTasks(ctx, func(_ context.Context, task *contract.Task) error { got <- task return nil }) if err != nil { t.Fatalf("consume: %v", err) } - defer stop() + defer drain(context.Background()) // --- Gateway 发布一个任务 --- want := &contract.Task{ @@ -104,7 +104,7 @@ func TestToolCallRoundTrip(t *testing.T) { } defer srv.Close() - unsub, err := srv.ServeTool(contract.SubjectToolsGoAll, contract.QueueToolsGo, + drain, err := srv.ServeTool(contract.SubjectToolsGoAll, contract.QueueToolsGo, func(_ context.Context, call *contract.ToolCall) *contract.ToolResult { if call.Tool != "wiki_search" { return &contract.ToolResult{OK: false, Error: "unknown tool"} @@ -114,7 +114,7 @@ func TestToolCallRoundTrip(t *testing.T) { if err != nil { t.Fatalf("serve tool: %v", err) } - defer func() { _ = unsub() }() + defer drain(context.Background()) // --- Dispatcher 侧:同步调用工具 --- dp, err := bus.Connect(url) @@ -220,7 +220,7 @@ func TestConcurrentConsume(t *testing.T) { blockA := make(chan struct{}) doneB := make(chan string, 1) - stop, err := dp.ConsumeTasks(ctx, func(_ context.Context, task *contract.Task) error { + drain, err := dp.ConsumeTasks(ctx, func(_ context.Context, task *contract.Task) error { switch task.ID { case "A": <-blockA // 模拟审批长阻塞 @@ -232,7 +232,7 @@ func TestConcurrentConsume(t *testing.T) { if err != nil { t.Fatalf("consume: %v", err) } - defer stop() + defer drain(context.Background()) defer close(blockA) // 先发 A 并等它被消费、卡在 handler 里;再发 B。 @@ -267,7 +267,7 @@ func TestConcurrentToolServe(t *testing.T) { defer srv.Close() blockSlow := make(chan struct{}) - unsub, err := srv.ServeTool(contract.SubjectToolsGoAll, contract.QueueToolsGo, + drain, err := srv.ServeTool(contract.SubjectToolsGoAll, contract.QueueToolsGo, func(_ context.Context, call *contract.ToolCall) *contract.ToolResult { if call.Tool == "slow" { <-blockSlow // 模拟慢工具(如 RAG 嵌入) @@ -277,7 +277,7 @@ func TestConcurrentToolServe(t *testing.T) { if err != nil { t.Fatalf("serve tool: %v", err) } - defer unsub() + defer drain(context.Background()) defer close(blockSlow) cli, err := bus.Connect(url) @@ -306,3 +306,57 @@ func TestConcurrentToolServe(t *testing.T) { } t.Log("✓ 单实例工具并发生效:slow 仍阻塞时 fast 已返回") } + +// TestGracefulDrain 验证优雅停机:drain 会等在途任务跑完(不被信号掐断), +// 且在途任务的 ctx 不随消费停止而取消。 +func TestGracefulDrain(t *testing.T) { + url := startEmbeddedNATS(t) + gw, err := bus.Connect(url) + if err != nil { + t.Fatalf("gateway connect: %v", err) + } + defer gw.Close() + if err := gw.EnsureTaskStream(context.Background()); err != nil { + t.Fatalf("ensure stream: %v", err) + } + dp, err := bus.Connect(url) + if err != nil { + t.Fatalf("dispatcher connect: %v", err) + } + defer dp.Close() + + started := make(chan struct{}) + finished := make(chan struct{}) + drain, err := dp.ConsumeTasks(context.Background(), func(c context.Context, _ *contract.Task) error { + close(started) + time.Sleep(500 * time.Millisecond) // 模拟在途长任务 + if c.Err() != nil { // 关键:消费停止不应取消在途任务的 ctx + t.Errorf("在途任务 ctx 不应被取消: %v", c.Err()) + } + close(finished) + return nil + }) + if err != nil { + t.Fatalf("consume: %v", err) + } + if _, err := gw.PublishTask(context.Background(), &contract.Task{ID: "drain1", Graph: json.RawMessage(`{}`)}); err != nil { + t.Fatalf("publish: %v", err) + } + <-started // 等任务进入 handler + + t0 := time.Now() + dctx, cancel := context.WithTimeout(context.Background(), 5*time.Second) + defer cancel() + drain(dctx) // 应阻塞至在途任务跑完 + elapsed := time.Since(t0) + + select { + case <-finished: + default: + t.Fatal("drain 返回时在途任务竟未完成(被掐断或未等待)") + } + if elapsed < 300*time.Millisecond { + t.Fatalf("drain 过早返回(%v),未真正等待在途任务", elapsed) + } + t.Logf("✓ drain 等待在途任务跑完:耗时 %v", elapsed) +}