71102d2424
submitted → running → done / failed / timeout 显式状态机,根治"卡运行中 看不出来"(此前 Task.Status 只写死 submitted、从不流转)。 - contract:新增 SubjectTaskStatus 主题 + TaskStatusEvent + 状态常量 - shared/bus:PublishTaskStatus / SubscribeTaskStatus(core NATS pub-sub) - dispatcher:StatusSink 接口 + Orchestrator.Handle 状态钩子——进入执行 →running、收尾 finishStatus→done/failed、整体超时上限 taskExecTimeout =3min→timeout;经 sub 回写 - gateway:SubscribeTaskStatus 落 PG(Task 增 Detail 字段,AutoMigrate); 新增 GET /api/v1/tasks/:id 供 UI 轮询状态 验收:实测 submitted→running→done 流转 + PG 持久化 + 端点查询闭环; make test-go 全绿。HITL 中断恢复 / 多智能体仍按场景待做。 Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
90 lines
3.2 KiB
Go
90 lines
3.2 KiB
Go
// Package nats 是网关对共享 bus 的薄封装(发布任务 / 订阅 Token 回流)。
|
|
package nats
|
|
|
|
import (
|
|
"context"
|
|
"encoding/json"
|
|
"log"
|
|
|
|
sharedbus "github.com/sundynix/sundynix-shared/bus"
|
|
"github.com/sundynix/sundynix-shared/contract"
|
|
)
|
|
|
|
// Bus 包装共享 bus,向网关其余代码暴露发布能力。
|
|
type Bus struct {
|
|
inner *sharedbus.Bus
|
|
}
|
|
|
|
// MustConnect 接入 NATS 并确保任务流存在。
|
|
func MustConnect(url string) *Bus {
|
|
inner, err := sharedbus.Connect(url)
|
|
if err != nil {
|
|
log.Fatalf("[nats] connect: %v", err)
|
|
}
|
|
if err := inner.EnsureTaskStream(context.Background()); err != nil {
|
|
log.Fatalf("[nats] ensure stream: %v", err)
|
|
}
|
|
log.Printf("[nats] connected %s, task stream ready", url)
|
|
return &Bus{inner: inner}
|
|
}
|
|
|
|
// PublishTask 把组装后的 Task 发布到 sundynix.tasks.<id>。
|
|
func (b *Bus) PublishTask(ctx context.Context, t *contract.Task) error {
|
|
seq, err := b.inner.PublishTask(ctx, t)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
log.Printf("[nats] published task %s (seq=%d)", t.ID, seq)
|
|
return nil
|
|
}
|
|
|
|
// SubscribeTokens 订阅 sundynix.streams.<taskID> 的 Token 回流,
|
|
// 每个 Token 触发 onToken,流结束触发 onDone,返回 unsub。
|
|
func (b *Bus) SubscribeTokens(taskID string, onToken func([]byte), onDone func()) (func() error, error) {
|
|
return b.inner.SubscribeTokens(taskID, onToken, onDone)
|
|
}
|
|
|
|
// SubscribeExec 订阅 sundynix.exec.<taskID> 的执行轨迹事件(用于"运行·观测"SSE)。
|
|
func (b *Bus) SubscribeExec(taskID string, onEvent func([]byte), onDone func()) (func() error, error) {
|
|
return b.inner.SubscribeExec(taskID, onEvent, onDone)
|
|
}
|
|
|
|
// CallTool 经 NATS 同步调用一个 MCP 工具(用于网关侧写偏好记忆等)。
|
|
func (b *Bus) CallTool(ctx context.Context, subject string, call *contract.ToolCall) (*contract.ToolResult, error) {
|
|
return b.inner.CallTool(ctx, subject, call)
|
|
}
|
|
|
|
// Ping 同步探测某节点健康(如 dispatcher 心跳主题)。无人应答 / 超时即返回错误(视为下线)。
|
|
func (b *Bus) Ping(ctx context.Context, subject string) ([]byte, error) {
|
|
return b.inner.Ping(ctx, subject)
|
|
}
|
|
|
|
// SubscribeTaskStatus 订阅 dispatcher 回写的任务生命周期状态(落 PG)。
|
|
func (b *Bus) SubscribeTaskStatus(onEvent func(*contract.TaskStatusEvent)) (func() error, error) {
|
|
return b.inner.SubscribeTaskStatus(onEvent)
|
|
}
|
|
|
|
// ServeConfig 让网关作为配置控制面,响应某 kind 的配置请求。
|
|
func (b *Bus) ServeConfig(kind string, provide func() *contract.ModelConfig) (func() error, error) {
|
|
return b.inner.ServeConfig(kind, provide)
|
|
}
|
|
|
|
// PublishConfigUpdated 广播某 kind 的配置变更。
|
|
func (b *Bus) PublishConfigUpdated(kind string, cfg *contract.ModelConfig) error {
|
|
return b.inner.PublishConfigUpdated(kind, cfg)
|
|
}
|
|
|
|
// PublishIngest 把一条入库进度事件发到 sundynix.streams.<jobID>。
|
|
func (b *Bus) PublishIngest(jobID string, ev *contract.IngestEvent) error {
|
|
data, err := json.Marshal(ev)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
return b.inner.PublishToken(jobID, data)
|
|
}
|
|
|
|
// CompleteStream 发送入库流结束信号。
|
|
func (b *Bus) CompleteStream(jobID string) error { return b.inner.CompleteStream(jobID) }
|
|
|
|
func (b *Bus) Close() { b.inner.Close() }
|