Files
sundynix-agentix/sundynix-gateway/internal/nats/publisher.go
T
Blizzard 10f08ffb14 feat(harness): 评测闭环 —— 评测结果落库 + 分级 + 告警 + 可查(测温计→恒温器)
此前 eval 只打日志、不闭环。现在:
- 分级:evalLevel 据综合分+忠实度 → ok(≥0.75) / warn(≥0.5 或忠实<0.6) / poor(<0.5);poor 出 slog.Warn 告警。
- 落库:dispatcher 评完经 NATS(SubjectEval) 广播 EvalEvent → 网关订阅写 PG(新表 sundynix_eval,
  按 task_id upsert)。沿用任务状态回写那套(dispatcher 无 DB,经 bus→gateway 落库)。
- 可查:GET /api/v1/tasks/:id/eval 返回 overall/rule/llm/faithful/level/flags/reason/sources。
- 契约 EvalEvent + EvalOK/Warn/Poor;bus PublishEval/SubscribeEval;dispatcher EvalSink(NewOrchestrator 第9参)。

验证:三模块 build+vet+test 全绿;live RAG 任务评测落库,端点返回 overall~1.0 / level=ok / faithful=1 / sources=1。
剩:桌面端质量面板、低分自动重试(P3)。project_analysis 勾掉该项。

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-06-25 15:26:19 +08:00

100 lines
3.6 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)
}
// PublishApproval 把一次人工审批决定发给 dispatcher(解除审批节点阻塞)。
func (b *Bus) PublishApproval(dec *contract.ApprovalDecision) error {
return b.inner.PublishApproval(dec)
}
// SubscribeEval 订阅 dispatcher 回写的自动化评测结果(落 PG)。
func (b *Bus) SubscribeEval(onEvent func(*contract.EvalEvent)) (func() error, error) {
return b.inner.SubscribeEval(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() }