c7a02c3905
5 层 + 1 条 NATS 零拷贝消息总线的 monorepo(Monolith First → Microservices Morph B)。 纵向主干(任务流 + Token 流回流)已真实跑通,横向各层能力为带注释的桩。 已贯通(real code): - sundynix-shared: 共享契约 + JetStream/core NATS 真实收发(bus) + 内嵌 NATS(devnats) + e2e 测试 - sundynix-gateway: Gin 接入 + DSL 解析组装 + NATS Publish + SSE 流式输出 - sundynix-dispatcher: NATS 消费 + Eino Orchestrator 流式回流 + 熔断器 + LLM Pool 占位流式 - 链路: HTTP POST → DSL → sundynix.tasks.* → Dispatcher → Token 经 sundynix.streams.<id> 回流 → SSE - 基础设施: docker-compose(nats/postgres/redis/neo4j/milvus) + Makefile(make demo/e2e) 待填(桩): - Eino 图编排 compose.NewGraph、LLM Pool 接 vLLM/Ollama - Gateway store 换真实 pgx/redis - sundynix-mcp-go: Bleve+Milvus+Neo4j 混合检索 / UniOffice / 外部 API - sundynix-mcp-py: gVisor 沙箱 / MinerU(PaddleOCR) / Docker 解释器 - sundynix-desktop: React Flow 画布 → DSL 导出 → SSE 展示
150 lines
3.9 KiB
Go
150 lines
3.9 KiB
Go
package bus_test
|
||
|
||
import (
|
||
"context"
|
||
"encoding/json"
|
||
"testing"
|
||
"time"
|
||
|
||
natsserver "github.com/nats-io/nats-server/v2/server"
|
||
natstest "github.com/nats-io/nats-server/v2/test"
|
||
|
||
"github.com/sundynix/sundynix-shared/bus"
|
||
"github.com/sundynix/sundynix-shared/contract"
|
||
)
|
||
|
||
// startEmbeddedNATS 启动一个内嵌、开启 JetStream 的 NATS 服务器,免 Docker。
|
||
func startEmbeddedNATS(t *testing.T) string {
|
||
t.Helper()
|
||
opts := natstest.DefaultTestOptions
|
||
opts.Port = -1 // 随机端口
|
||
opts.JetStream = true
|
||
opts.StoreDir = t.TempDir()
|
||
srv := natstest.RunServer(&opts)
|
||
if !srv.ReadyForConnections(5 * time.Second) {
|
||
t.Fatal("embedded NATS not ready")
|
||
}
|
||
t.Cleanup(srv.Shutdown)
|
||
_ = natsserver.Server{} // 触发包引用
|
||
return srv.ClientURL()
|
||
}
|
||
|
||
// TestTaskRoundTrip 模拟 Gateway 发布 → NATS → Dispatcher 消费 的完整任务流。
|
||
func TestTaskRoundTrip(t *testing.T) {
|
||
url := startEmbeddedNATS(t)
|
||
|
||
// --- Gateway 侧:连接并声明任务流 ---
|
||
gw, err := bus.Connect(url)
|
||
if err != nil {
|
||
t.Fatalf("gateway connect: %v", err)
|
||
}
|
||
defer gw.Close()
|
||
|
||
ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
|
||
defer cancel()
|
||
|
||
if err := gw.EnsureTaskStream(ctx); err != nil {
|
||
t.Fatalf("ensure stream: %v", err)
|
||
}
|
||
|
||
// --- Dispatcher 侧:连接并开始消费 ---
|
||
dp, err := bus.Connect(url)
|
||
if err != nil {
|
||
t.Fatalf("dispatcher connect: %v", err)
|
||
}
|
||
defer dp.Close()
|
||
|
||
got := make(chan *contract.Task, 1)
|
||
stop, 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()
|
||
|
||
// --- Gateway 发布一个任务 ---
|
||
want := &contract.Task{
|
||
ID: "task_demo_001",
|
||
Graph: json.RawMessage(`{"nodes":[{"id":"n1","type":"agent"}],"edges":[]}`),
|
||
Meta: map[string]any{"user": "wt"},
|
||
}
|
||
seq, err := gw.PublishTask(ctx, want)
|
||
if err != nil {
|
||
t.Fatalf("publish: %v", err)
|
||
}
|
||
if seq == 0 {
|
||
t.Fatal("expected non-zero stream sequence")
|
||
}
|
||
|
||
// --- 断言 Dispatcher 收到同一个任务 ---
|
||
select {
|
||
case task := <-got:
|
||
if task.ID != want.ID {
|
||
t.Fatalf("task id = %q, want %q", task.ID, want.ID)
|
||
}
|
||
if task.Meta["user"] != "wt" {
|
||
t.Fatalf("task meta lost: %+v", task.Meta)
|
||
}
|
||
t.Logf("✓ 任务流打通:Gateway publish (seq=%d) → NATS → Dispatcher consume,task_id=%s", seq, task.ID)
|
||
case <-time.After(5 * time.Second):
|
||
t.Fatal("timeout: dispatcher 未收到任务")
|
||
}
|
||
}
|
||
|
||
// TestTokenStreamRoundTrip 模拟 Dispatcher 回流 Token → Gateway 订阅 的流式闭环。
|
||
func TestTokenStreamRoundTrip(t *testing.T) {
|
||
url := startEmbeddedNATS(t)
|
||
|
||
// Gateway 侧:先订阅(core NATS 无持久化,须先连)。
|
||
gw, err := bus.Connect(url)
|
||
if err != nil {
|
||
t.Fatalf("gateway connect: %v", err)
|
||
}
|
||
defer gw.Close()
|
||
|
||
const taskID = "task_stream_001"
|
||
var got []string
|
||
done := make(chan struct{})
|
||
unsub, err := gw.SubscribeTokens(taskID,
|
||
func(tok []byte) { got = append(got, string(tok)) },
|
||
func() { close(done) },
|
||
)
|
||
if err != nil {
|
||
t.Fatalf("subscribe tokens: %v", err)
|
||
}
|
||
defer func() { _ = unsub() }()
|
||
|
||
// Dispatcher 侧:逐 Token 回流后发结束信号。
|
||
dp, err := bus.Connect(url)
|
||
if err != nil {
|
||
t.Fatalf("dispatcher connect: %v", err)
|
||
}
|
||
defer dp.Close()
|
||
|
||
want := []string{"Hello", " ", "Agent", "!"}
|
||
for _, tok := range want {
|
||
if err := dp.PublishToken(taskID, []byte(tok)); err != nil {
|
||
t.Fatalf("publish token: %v", err)
|
||
}
|
||
}
|
||
if err := dp.CompleteStream(taskID); err != nil {
|
||
t.Fatalf("complete stream: %v", err)
|
||
}
|
||
|
||
select {
|
||
case <-done:
|
||
joined := ""
|
||
for _, s := range got {
|
||
joined += s
|
||
}
|
||
if joined != "Hello Agent!" {
|
||
t.Fatalf("token stream = %q, want %q", joined, "Hello Agent!")
|
||
}
|
||
t.Logf("✓ Token 流闭环:Dispatcher 回流 %d 个 token → Gateway 拼回 %q", len(got), joined)
|
||
case <-time.After(5 * time.Second):
|
||
t.Fatal("timeout: 未收到流结束信号")
|
||
}
|
||
}
|