Files
sundynix-agentix/sundynix-gateway/internal/nats/publisher.go
T
Blizzard 9b153871eb feat(monitor): NATS 集群 Raft 副本健康 + 基建 ping 延迟(监测组完善)
此前 /status 把 NATS 当一盏二元灯(连不上就 fatal 故恒真),看不出 3 节点集群里
哪个节点掉了、JetStream 持久流的 Raft 副本是否还齐(计费/状态/评测流不丢的关键)。
DB/Redis/MinIO 也只二元 ping、无延迟。

- bus.ClusterStatus:连的节点名 + 集群发现节点数(nc.Servers) + RTT(nc.RTT) + 6 条关键
  持久流(tasks/status/usage/eval/ingest/approvals)的 Raft 副本健康(leader + healthy/total,
  单节点部署记 1/1;某节点掉队 → healthy<total 记降级)。
- /status 新增 nats 集群对象 + NATS 灯改为'连接且无副本降级才绿'、detail 显示'N 节点·连 X·
  流副本齐全/降级'、latency=RTT;PG/Redis/MinIO 加 ping 往返耗时。
- admin StatusPage:新增 NATS 集群面板(节点/RTT/连接 + 各流 leader/副本健康/消息数表),
  基建行显示延迟。

sundynix-shared/gateway build+vet+test 绿;admin tsc 绿。

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-07-21 16:15:22 +08:00

149 lines
5.8 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)
}
// 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() }