diff --git a/sundynix-shared/bus/bus.go b/sundynix-shared/bus/bus.go index 7f08bd9..a33b201 100644 --- a/sundynix-shared/bus/bus.go +++ b/sundynix-shared/bus/bus.go @@ -6,6 +6,9 @@ import ( "context" "encoding/json" "fmt" + "log" + "os" + "strconv" "time" "github.com/nats-io/nats.go" @@ -441,9 +444,21 @@ func (b *Bus) SubscribeConfigUpdated(kind string, onUpdate func(*contract.ModelC // TaskHandler 处理一个消费到的任务。 type TaskHandler func(ctx context.Context, t *contract.Task) error +// taskConcurrency 返回任务并发处理上限(env DISPATCHER_CONCURRENCY,默认 8)。 +func taskConcurrency() int { + if v := os.Getenv("DISPATCHER_CONCURRENCY"); v != "" { + if n, err := strconv.Atoi(v); err == nil && n > 0 { + return n + } + } + return 8 +} + // ConsumeTasks 在持久消费者上消费任务,队列组内负载均衡。 -// 返回的 stop 函数用于优雅停止消费。 +// 每个任务分发到独立 worker goroutine 并发执行——一个慢任务/HITL 待审不再阻塞后续任务。 +// 并发上限由信号量 + 消费者 MaxAckPending 双重约束(背压)。返回的 stop 用于优雅停止消费。 func (b *Bus) ConsumeTasks(ctx context.Context, h TaskHandler) (stop func(), err error) { + concurrency := taskConcurrency() cons, err := b.js.CreateOrUpdateConsumer(ctx, contract.StreamTasks, jetstream.ConsumerConfig{ Durable: contract.ConsumerDurable, AckPolicy: jetstream.AckExplicitPolicy, @@ -451,31 +466,47 @@ func (b *Bus) ConsumeTasks(ctx context.Context, h TaskHandler) (stop func(), err // HITL:审批节点会让 Handle 阻塞等人工决定(最长约 5 分钟), // AckWait 必须覆盖「审批等待 + 图执行」总时长,否则消息在途未 ack 会被重投成重复任务。 AckWait: 15 * time.Minute, + // 在途未 ack 上限 = 并发上限:服务端不会下发超过本节点同时能处理的量(背压)。 + MaxAckPending: concurrency, }) if err != nil { return nil, fmt.Errorf("create consumer: %w", err) } + sem := make(chan struct{}, concurrency) // 限并发:最多 N 个任务同时执行 cc, err := cons.Consume(func(msg jetstream.Msg) { t, err := contract.Unmarshal(msg.Data()) if err != nil { _ = msg.Term() // 脏数据,丢弃不重投 return } - // 从消息头还原上游链路,开消费 span(成为 gateway 发布 span 的子节点)。 - mctx := extractTrace(ctx, 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))) - herr := h(mctx, t) - if herr != nil { - span.RecordError(herr) - } - span.End() - if herr != nil { - _ = msg.NakWithDelay(time.Second) // 处理失败,延迟重投 + // 背压:并发已满则在此等空位;正在关停则留消息不 ack(稍后重投)。 + select { + case sem <- struct{}{}: + case <-ctx.Done(): return } - _ = msg.Ack() + go func() { + defer func() { + <-sem // 释放并发额度 + 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())) + mctx, span := tracer().Start(mctx, "nats.consume task", + trace.WithSpanKind(trace.SpanKindConsumer), + trace.WithAttributes(attribute.String("sundynix.task_id", t.ID))) + defer span.End() + if herr := h(mctx, t); herr != nil { + span.RecordError(herr) + _ = msg.NakWithDelay(time.Second) // 处理失败,延迟重投 + return + } + _ = msg.Ack() + }() }) if err != nil { return nil, fmt.Errorf("consume: %w", err) diff --git a/sundynix-shared/bus/bus_e2e_test.go b/sundynix-shared/bus/bus_e2e_test.go index bed00f8..80cbcbe 100644 --- a/sundynix-shared/bus/bus_e2e_test.go +++ b/sundynix-shared/bus/bus_e2e_test.go @@ -194,3 +194,63 @@ func TestTokenStreamRoundTrip(t *testing.T) { t.Fatal("timeout: 未收到流结束信号") } } + +// TestConcurrentConsume 验证并发消费:一个长时间阻塞的任务(模拟 HITL 待审) +// 不应阻塞后续任务——旧的串行消费下本测试会超时失败。 +func TestConcurrentConsume(t *testing.T) { + url := startEmbeddedNATS(t) + + gw, err := bus.Connect(url) + if err != nil { + t.Fatalf("gateway connect: %v", err) + } + defer gw.Close() + + ctx, cancel := context.WithTimeout(context.Background(), 15*time.Second) + defer cancel() + if err := gw.EnsureTaskStream(ctx); 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() + + blockA := make(chan struct{}) + doneB := make(chan string, 1) + stop, err := dp.ConsumeTasks(ctx, func(_ context.Context, task *contract.Task) error { + switch task.ID { + case "A": + <-blockA // 模拟审批长阻塞 + case "B": + doneB <- task.ID + } + return nil + }) + if err != nil { + t.Fatalf("consume: %v", err) + } + defer stop() + defer close(blockA) + + // 先发 A 并等它被消费、卡在 handler 里;再发 B。 + if _, err := gw.PublishTask(ctx, &contract.Task{ID: "A", Graph: json.RawMessage(`{}`)}); err != nil { + t.Fatalf("publish A: %v", err) + } + time.Sleep(400 * time.Millisecond) + if _, err := gw.PublishTask(ctx, &contract.Task{ID: "B", Graph: json.RawMessage(`{}`)}); err != nil { + t.Fatalf("publish B: %v", err) + } + + select { + case id := <-doneB: + if id != "B" { + t.Fatalf("unexpected task done: %q", id) + } + t.Log("✓ 并发消费生效:A 仍阻塞时 B 已完成") + case <-time.After(4 * time.Second): + t.Fatal("B 被阻塞的 A 卡住——并发消费未生效(仍是串行)") + } +}