Files
sundynix-agentix/sundynix-gateway/internal/nats/publisher.go
T

165 lines
6.6 KiB
Go
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
// 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)
}
if err := inner.EnsureIngestStream(context.Background()); err != nil {
log.Fatalf("[nats] ensure ingest stream: %v", err)
}
if err := inner.EnsureStatusStream(context.Background()); err != nil {
log.Fatalf("[nats] ensure status stream: %v", err)
}
if err := inner.EnsureUsageStream(context.Background()); err != nil {
log.Fatalf("[nats] ensure usage stream: %v", err)
}
if err := inner.EnsureEvalStream(context.Background()); err != nil {
log.Fatalf("[nats] ensure eval stream: %v", err)
}
log.Printf("[nats] connected %s, task + ingest + status + usage + eval streams ready", url)
return &Bus{inner: inner}
}
// ClusterStatus 透传共享 bus 的 NATS 集群体检(供「服务状态」监测面板)。
func (b *Bus) ClusterStatus(ctx context.Context) sharedbus.NATSClusterStatus {
return b.inner.ClusterStatus(ctx)
}
// 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)
}
// ServeTool 以队列组订阅一族工具主题并按调用分发(平台工具族:gateway 自己当工具提供方,
// 见 JARVIS_BRAIN_DESIGN.md §2.1)。返回 drain 供优雅停机。
func (b *Bus) ServeTool(subject, queue string, h sharedbus.ToolHandler) (func(context.Context), error) {
return b.inner.ServeTool(subject, queue, h)
}
// PublishVoiceEvent 给某用户的语音会话发事件(JARVIS 动作/主动播报;无会话则静默丢弃)。
func (b *Bus) PublishVoiceEvent(uid string, ev *contract.VoiceEvent) error {
return b.inner.PublishVoiceEvent(uid, ev)
}
// SubscribeVoiceEvent 订阅某用户的语音事件(语音 WS 会话生命周期内挂载)。
func (b *Bus) SubscribeVoiceEvent(uid string, onEvent func(*contract.VoiceEvent)) (func() error, error) {
return b.inner.SubscribeVoiceEvent(uid, onEvent)
}
// Ping 同步探测某节点健康(如 dispatcher 心跳主题)。无人应答 / 超时即返回错误(视为下线)。
func (b *Bus) Ping(ctx context.Context, subject string) ([]byte, error) {
return b.inner.Ping(ctx, subject)
}
// ConsumeTaskStatus 持久消费 dispatcher 回写的任务生命周期状态(落 PGat-least-once + 幂等)。
func (b *Bus) ConsumeTaskStatus(ctx context.Context, h func(context.Context, *contract.TaskStatusEvent) error) (func(context.Context), error) {
return b.inner.ConsumeTaskStatus(ctx, h)
}
// PublishApproval 把一次人工审批决定发给 dispatcher(解除审批节点阻塞)。
func (b *Bus) PublishApproval(dec *contract.ApprovalDecision) error {
return b.inner.PublishApproval(dec)
}
// ConsumeEval 持久消费 dispatcher 回写的自动化评测结果(落 PGat-least-once + 幂等)。
func (b *Bus) ConsumeEval(ctx context.Context, h func(context.Context, *contract.EvalEvent) error) (func(context.Context), error) {
return b.inner.ConsumeEval(ctx, h)
}
// ConsumeUsage 持久消费 dispatcher 回写的任务 token 用量(计费,at-least-once + 幂等)。
func (b *Bus) ConsumeUsage(ctx context.Context, h func(context.Context, *contract.UsageEvent) error) (func(context.Context), error) {
return b.inner.ConsumeUsage(ctx, h)
}
// 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) }
// ---- 入库工作队列(JetStream 持久)----
// EnsureIngestStream 确保入库作业流存在(启动时调)。
func (b *Bus) EnsureIngestStream(ctx context.Context) error { return b.inner.EnsureIngestStream(ctx) }
// PublishIngestJob 把入库作业持久入队(崩溃不丢,由 worker 消费)。
func (b *Bus) PublishIngestJob(ctx context.Context, job *contract.IngestJob) error {
return b.inner.PublishIngestJob(ctx, job)
}
// ConsumeIngestJobs 启动入库 worker 池消费作业(有界并发 + 背压 + 崩溃重投)。
func (b *Bus) ConsumeIngestJobs(ctx context.Context, h sharedbus.IngestHandler) (func(context.Context), error) {
return b.inner.ConsumeIngestJobs(ctx, h)
}
// ---- Prompt 控制面 ----
// ServePrompts 让网关响应「取全部激活 prompt」请求。
func (b *Bus) ServePrompts(provide func() map[string]string) (func() error, error) {
return b.inner.ServePrompts(provide)
}
// PublishPromptsUpdated 广播最新激活集 → 各服务热更新。
func (b *Bus) PublishPromptsUpdated(m map[string]string) error {
return b.inner.PublishPromptsUpdated(m)
}
func (b *Bus) Close() { b.inner.Close() }