From 16a6c4b1aad35eca60aa3d73f0c4950e639638a9 Mon Sep 17 00:00:00 2001 From: Blizzard Date: Wed, 24 Jun 2026 13:18:47 +0800 Subject: [PATCH] =?UTF-8?q?feat(hitl):=20=E4=BA=BA=E5=B7=A5=E5=AE=A1?= =?UTF-8?q?=E6=89=B9=E4=B8=AD=E6=96=AD=EF=BC=88Eino=20Phase=20D=EF=BC=89?= =?UTF-8?q?=E2=80=94=E2=80=94=20=E5=AE=A1=E6=89=B9=E8=8A=82=E7=82=B9?= =?UTF-8?q?=E6=9A=82=E5=81=9C=E2=86=92=E6=89=B9=E5=87=86=E7=BB=AD=E8=B7=91?= =?UTF-8?q?/=E6=8B=92=E7=BB=9D=E4=B8=AD=E6=AD=A2?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 专用「审批」节点方案,全栈打通。 后端: - contract:TaskWaiting/TaskRejected 状态 + ApprovalSubject/ApprovalDecision。 - bus:PublishApproval + WaitApproval(订阅决定主题,带超时);消费者 AckWait 30s→15min(阻塞等人审期间消息未 ack,否则重投成重复任务)。 - dispatcher:ApprovalWaiter 接口 + approvalNode——执行到审批节点发 await 事件 + 置 waiting,阻塞等决定。批准→回 running 放行下游;拒绝/超时→errRejected 哨兵→ 剪下游→rejected,优雅收尾不计熔断。超时安全默认拒绝。 taskExecTimeout 3→10min(审批5 < 执行10 < AckWait15)。 - 网关:POST /tasks/:id/approve(仅 waiting 态受理,幂等)。 桌面端: - Studio 新增「人工审批」节点(nodeCatalog,可填标题/说明)。 - run.ts pendingApproval() 从 exec 流派生待审中断 + waiting 节点状态。 - BottomDrawer ApprovalBar:琥珀审批条(摘要 + 批准/拒绝 + 备注,调 api.approveTask)。 - ExecTrace waiting 状态灯。 验证:后端 curl 实测 waiting→批准→running→done;waiting→拒绝→rejected(下游未跑)。 前端 tsc + vitest 36 过(pendingApproval 4 例 + waiting 状态)。全模块 build+vet+test 全绿。 EINO_ADOPTION Phase D 标记 HITL 完成。 未覆盖:compose 路径(EINO_COMPOSE,默认关)的 approval 节点;桌面端审批条未在 GUI 实点。 Co-Authored-By: Claude Opus 4.8 (1M context) --- EINO_ADOPTION.md | 2 +- .../frontend/src/components/ExecTrace.tsx | 1 + sundynix-desktop/frontend/src/lib/api.ts | 12 ++++ sundynix-desktop/frontend/src/lib/run.test.ts | 34 +++++++++- sundynix-desktop/frontend/src/lib/run.ts | 27 +++++++- .../frontend/src/shell/BottomDrawer.tsx | 62 ++++++++++++++++++- .../frontend/src/studio/nodeCatalog.ts | 12 ++++ sundynix-dispatcher/cmd/dispatcher/main.go | 5 +- sundynix-dispatcher/internal/eino/graph.go | 51 +++++++++++++++ .../internal/eino/orchestrator.go | 47 ++++++++++---- .../internal/nats/subscriber.go | 5 ++ .../internal/handler/task_handler.go | 28 +++++++++ sundynix-gateway/internal/nats/publisher.go | 5 ++ sundynix-gateway/internal/router/router.go | 39 ++++++------ sundynix-shared/bus/bus.go | 45 ++++++++++++++ sundynix-shared/contract/task.go | 24 ++++++- 16 files changed, 360 insertions(+), 39 deletions(-) diff --git a/EINO_ADOPTION.md b/EINO_ADOPTION.md index 1ac2753..b3774c8 100644 --- a/EINO_ADOPTION.md +++ b/EINO_ADOPTION.md @@ -10,7 +10,7 @@ - [x] **Phase A · 地基**:`llm.Pool` 换 Eino ChatModel 组件(commit d84b1ec,验收通过) - [x] **Phase B · 质变**:MCP 工具→`InvokableTool` + ReAct agent(模型自主调工具,验收 7/7 命中) - [x] **Phase C · 编排归一**:✅ 全图 `DSL→compose.Graph` 编译器(全节点 + branch + DAG 并行调度)+ callbacks→ExecEvent 归一,等价回归通过(EINO_COMPOSE 灰度开关,默认关;过渡期后 graph.go 退役) -- [~] **Phase D · 状态化执行**:✅ 任务生命周期 FSM(已完成)/ ⬜ HITL 中断恢复 / ⬜ 多智能体(按场景) +- [~] **Phase D · 状态化执行**:✅ 任务生命周期 FSM / ✅ HITL 人工审批中断(审批节点暂停→waiting→批准回 running / 拒绝→rejected,NATS 决定回传 + AckWait 续租,全栈含 Studio 审批节点 + 运行抽屉批准条)/ ⬜ 多智能体(按场景) **组件化补完(A)**:检索 → `ragRetriever`(`components/retriever.Retriever`,`eino_components.go`);提示词 → `buildMessages` 改用 `prompt.FromMessages`+`MessagesPlaceholder`;工具 → `mcpTool`(`InvokableTool`,模型自主调用)。至此终态架构 8 层中 模型/工具/检索/提示词/编排/智能体(单)/可观测 均已 Eino 组件化;剩 人机交互(中断恢复,按场景)。 **自主 agent 工具集动态化**:agent 工具集不再硬编码——mcp-go 注册表(单一事实源)每个工具声明 `agent/params/inject`,`list_tools` 上报,dispatcher `agentTools()` 动态发现并建 `InvokableTool`、运行时注入 `user_id/session_id/kb`(不暴露给模型)。**加工具只改 mcp-go 注册表,dispatcher 零改动。** 当前暴露 4 个(wiki_search/recall_user_memory/remember_user_fact/history_get);实测模型自主调用新暴露的 remember_user_fact 成功(参数自生成、user_id 服务端注入)。 diff --git a/sundynix-desktop/frontend/src/components/ExecTrace.tsx b/sundynix-desktop/frontend/src/components/ExecTrace.tsx index a07eebb..dfc6968 100644 --- a/sundynix-desktop/frontend/src/components/ExecTrace.tsx +++ b/sundynix-desktop/frontend/src/components/ExecTrace.tsx @@ -37,6 +37,7 @@ function StatusDot({ status }: { status: NodeTrace["status"] }) { if (status === "running") return ; if (status === "done") return ; if (status === "error") return ; + if (status === "waiting") return ; // 待审批 return ; } diff --git a/sundynix-desktop/frontend/src/lib/api.ts b/sundynix-desktop/frontend/src/lib/api.ts index ff83de4..ad0e9e6 100644 --- a/sundynix-desktop/frontend/src/lib/api.ts +++ b/sundynix-desktop/frontend/src/lib/api.ts @@ -109,6 +109,18 @@ export async function submitTask(dsl: TaskDsl, id: Identity): Promise { return data.task_id; } +// approveTask: POST /api/v1/tasks/:id/approve —— HITL 人工审批决定(批准放行 / 拒绝中止)。 +export async function approveTask(taskId: string, approved: boolean, opts?: { node?: string; note?: string }): Promise { + const res = guard401( + await fetch(`${GATEWAY}/api/v1/tasks/${taskId}/approve`, { + method: "POST", + headers: { "Content-Type": "application/json", ...bearer() }, + body: JSON.stringify({ approved, node: opts?.node ?? "", note: opts?.note ?? "" }), + }), + ); + if (!res.ok) throw new Error(`approve failed: ${res.status} ${await res.text()}`); +} + // streamTokens: 订阅 SSE /api/v1/tasks/:id/stream,逐 token 回调,done 收尾。 // 返回关闭函数。注意 EventSource 无法带请求头,但流按 task_id 寻址,无需身份头。 export function streamTokens( diff --git a/sundynix-desktop/frontend/src/lib/run.test.ts b/sundynix-desktop/frontend/src/lib/run.test.ts index 286a78d..5efe026 100644 --- a/sundynix-desktop/frontend/src/lib/run.test.ts +++ b/sundynix-desktop/frontend/src/lib/run.test.ts @@ -1,6 +1,6 @@ import { describe, it, expect } from "vitest"; import type { ExecEvent } from "./api"; -import { deriveNodes } from "./run"; +import { deriveNodes, pendingApproval } from "./run"; let seq = 0; function ev(node: string, phase: string, extra: Partial = {}): ExecEvent { @@ -48,4 +48,36 @@ describe("deriveNodes(ExecEvent 流 → 节点轨迹)", () => { const [n] = deriveNodes([ev("a", "start", { label: "旧" }), ev("a", "end", { label: "新" })]); expect(n.label).toBe("新"); }); + + it("await 事件 → 节点状态 waiting", () => { + const [n] = deriveNodes([ev("ap", "await", { kind: "approval", label: "审批" })]); + expect(n.status).toBe("waiting"); + }); +}); + +describe("pendingApproval(HITL 待审批中断)", () => { + it("await 未收口 → 返回待审批", () => { + const p = pendingApproval([ev("ap", "await", { kind: "approval", label: "高危审批", detail: "摘要" })]); + expect(p).toEqual({ node: "ap", title: "高危审批", summary: "摘要" }); + }); + + it("await 后 end(同节点)→ 已决,返回 null", () => { + const p = pendingApproval([ + ev("ap", "await", { kind: "approval" }), + ev("ap", "end", { kind: "approval" }), + ]); + expect(p).toBeNull(); + }); + + it("await 后 error(超时/拒绝)→ 已决,返回 null", () => { + const p = pendingApproval([ + ev("ap", "await", { kind: "approval" }), + ev("ap", "error", { kind: "approval" }), + ]); + expect(p).toBeNull(); + }); + + it("无审批事件 → null", () => { + expect(pendingApproval([ev("a", "start"), ev("a", "end")])).toBeNull(); + }); }); diff --git a/sundynix-desktop/frontend/src/lib/run.ts b/sundynix-desktop/frontend/src/lib/run.ts index f22d5ca..fbb0211 100644 --- a/sundynix-desktop/frontend/src/lib/run.ts +++ b/sundynix-desktop/frontend/src/lib/run.ts @@ -21,7 +21,7 @@ export const emptyRun: RunState = { phase: "idle", output: "", events: [], exec: // ---- 执行轨迹派生:把扁平 ExecEvent 流归并为按节点聚合的轨迹 ---- -export type NodeStatus = "running" | "done" | "error" | "info"; +export type NodeStatus = "running" | "done" | "error" | "info" | "waiting"; export interface NodeTrace { node: string; @@ -34,6 +34,28 @@ export interface NodeTrace { order: number; } +// PendingApproval 是一个待人工审批的中断(HITL):审批节点已发 await 事件、尚未被 end/error 收口。 +export interface PendingApproval { + node: string; + title: string; + summary: string; +} + +// pendingApproval 从执行事件流里找出当前待审批的中断:取最后一个 kind=approval & phase=await, +// 且其后该节点没有 end/error(未被批准/拒绝收口)的事件。无则 null。 +export function pendingApproval(events: ExecEvent[]): PendingApproval | null { + let pending: PendingApproval | null = null; + for (const e of events) { + if (e.kind !== "approval") continue; + if (e.phase === "await") { + pending = { node: e.node, title: e.label || "人工审批", summary: e.detail || "" }; + } else if (e.phase === "end" || e.phase === "error") { + if (pending && pending.node === e.node) pending = null; // 已决,收口 + } + } + return pending; +} + // deriveNodes 把事件流按 node 归并:start→running,end→done(带耗时),error→error,info→点事件/附注。 export function deriveNodes(events: ExecEvent[]): NodeTrace[] { const map = new Map(); @@ -50,6 +72,9 @@ export function deriveNodes(events: ExecEvent[]): NodeTrace[] { case "start": if (n.status !== "done" && n.status !== "error") n.status = "running"; break; + case "await": // HITL 审批节点暂停,等人工决定 + n.status = "waiting"; + break; case "end": n.status = "done"; n.ms = e.ms; diff --git a/sundynix-desktop/frontend/src/shell/BottomDrawer.tsx b/sundynix-desktop/frontend/src/shell/BottomDrawer.tsx index af4b093..e79ddb7 100644 --- a/sundynix-desktop/frontend/src/shell/BottomDrawer.tsx +++ b/sundynix-desktop/frontend/src/shell/BottomDrawer.tsx @@ -1,7 +1,8 @@ import { useState } from "react"; -import { ChevronDown, ChevronUp, Wrench } from "lucide-react"; -import { deriveNodes, type RunState } from "../lib/run"; +import { ChevronDown, ChevronUp, Wrench, ShieldCheck, Check, X } from "lucide-react"; +import { deriveNodes, pendingApproval, type RunState } from "../lib/run"; import { ExecTrace } from "../components/ExecTrace"; +import { approveTask } from "../lib/api"; import { Tabs, Badge, cn, type TabDef } from "../ui"; type Tab = "output" | "trace" | "tools" | "cite" | "eval"; @@ -26,8 +27,11 @@ export function BottomDrawer({ run }: { run: RunState }) { const statusText = run.phase === "streaming" ? "流式中…" : run.phase === "done" ? "完成 ✓" : run.phase === "error" ? `✗ ${run.error ?? "出错"}` : run.phase === "submitting" ? "提交中…" : "就绪"; + const approval = run.taskId ? pendingApproval(run.exec) : null; + return (
+ {approval && run.taskId && }
(null); + const [err, setErr] = useState(""); + + const decide = async (approved: boolean) => { + setBusy(approved ? "approve" : "reject"); + setErr(""); + try { + await approveTask(taskId, approved, { node, note }); + } catch (e) { + setErr((e as Error).message); + setBusy(null); + } + }; + + return ( +
+
+ + 人工审批 + {title} + 等待决定 +
+ {summary &&
{summary}
} +
+ setNote(e.target.value)} + placeholder="备注(可选,拒绝原因等)" + className="min-w-0 flex-1 rounded border border-line bg-ink-950/60 px-2 py-1 text-[11px] text-slate-200 placeholder:text-slate-600 focus:border-amber-500/50 focus:outline-none" + /> + + +
+ {err &&

提交失败:{err}

} +
+ ); +} + // ToolCalls:从执行事件里筛出工具调用节点,逐条展示入参 → 产出 + 耗时/状态。 function ToolCalls({ run }: { run: RunState }) { const tools = deriveNodes(run.exec).filter((n) => n.kind === "tool"); diff --git a/sundynix-desktop/frontend/src/studio/nodeCatalog.ts b/sundynix-desktop/frontend/src/studio/nodeCatalog.ts index 800a988..e320146 100644 --- a/sundynix-desktop/frontend/src/studio/nodeCatalog.ts +++ b/sundynix-desktop/frontend/src/studio/nodeCatalog.ts @@ -102,6 +102,18 @@ export const NODE_KINDS: Record = { fields: [{ key: "condition", label: "条件", type: "text", placeholder: "score > 0.8" }], defaults: { condition: "" }, }, + approval: { + kind: "approval", + label: "人工审批", + accent: "border-l-amber-500", + badge: "bg-amber-100 text-amber-700", + desc: "HITL:暂停等人工批准", + fields: [ + { key: "title", label: "审批标题", type: "text", placeholder: "如:高危操作审批" }, + { key: "prompt", label: "审批说明", type: "textarea", placeholder: "向审批人说明待执行的操作…" }, + ], + defaults: { title: "人工审批", prompt: "请审批是否继续执行后续步骤" }, + }, map: { kind: "map", label: "并行 / Map", diff --git a/sundynix-dispatcher/cmd/dispatcher/main.go b/sundynix-dispatcher/cmd/dispatcher/main.go index b3d11d4..fd1655a 100644 --- a/sundynix-dispatcher/cmd/dispatcher/main.go +++ b/sundynix-dispatcher/cmd/dispatcher/main.go @@ -45,8 +45,9 @@ func main() { } go sub.FetchModelConfigWithRetry(context.Background(), pool.SetConfig) - // sub 同时作为 Token 回流(TokenSink)、MCP 工具调用(ToolCaller)、执行事件(ExecSink)与任务状态回写(StatusSink)出口。 - orch, err := eino.NewOrchestrator(pool, breaker, eval, sub, sub, sub, sub) + // sub 同时作为 Token 回流(TokenSink)、MCP 工具调用(ToolCaller)、执行事件(ExecSink)、 + // 任务状态回写(StatusSink)与 HITL 审批等待(ApprovalWaiter)出口。 + orch, err := eino.NewOrchestrator(pool, breaker, eval, sub, sub, sub, sub, sub) if err != nil { log.Fatalf("[dispatcher] build eino graph: %v", err) } diff --git a/sundynix-dispatcher/internal/eino/graph.go b/sundynix-dispatcher/internal/eino/graph.go index 39a09bc..3729e0f 100644 --- a/sundynix-dispatcher/internal/eino/graph.go +++ b/sundynix-dispatcher/internal/eino/graph.go @@ -32,6 +32,7 @@ type board struct { sections []reportSection // map 并行 fan-out 产出的分项成稿(供 render 多章渲染) answer string // 当前成稿(多 agent 协作时 = 最近一个 agent 的产出 = 成品) agentOut []string // 各上游 agent 的产出(按序),注入下游 agent 上下文以实现接力协作 + rejected bool // HITL 审批节点拒绝/超时 → 置位,runGraph 中止并返回 errRejected } // runGraph 按 DSL 图的真实拓扑与连线执行(替代旧的线性拍平 compileFlow)。 @@ -144,6 +145,8 @@ func (o *Orchestrator) runGraph(ctx context.Context, t *contract.Task, tr *execT o.renderNode(nctx, t.ID, n, b, tr) case "branch": propagate = o.branchNode(n, b, outE[n.ID], nodeByID, tr) + case "approval": + propagate = o.approvalNode(nctx, t.ID, n, b, tr, outE[n.ID]) case "map": o.mapNode(nctx, t.ID, n, b, tr) case "output": @@ -152,11 +155,18 @@ func (o *Orchestrator) runGraph(ctx context.Context, t *contract.Task, tr *execT tr.info(n.Kind+":"+n.ID, "system", labelOf(n, n.Kind), "未识别节点,跳过") } nspan.End() + if b.rejected { + break // 审批拒绝/超时:中止后续节点,整图判 rejected + } for _, tgt := range propagate { active[tgt] = true } } + if b.rejected { + return b.answer, errRejected // 合法终态,Handle 据此判 rejected 并优雅收尾 + } + // 图里无 agent 节点(纯工具/检索图)也要出一段模型答复,否则没有输出。 if b.answer == "" { o.runConversation(ctx, t.ID, b, plan.System, tr, "agent") @@ -369,6 +379,47 @@ func (o *Orchestrator) branchNode(n dsl.Node, b *board, outs []dsl.Edge, byID ma return chosen } +// approvalNode 是 HITL 人工审批中断:执行到此暂停,把待审摘要推给 UI(exec 事件 kind=approval/phase=await +// + 状态 waiting),阻塞等人工批准/拒绝(带超时,安全默认拒绝)。 +// 批准 → 状态回 running 并放行下游;拒绝/超时 → 置 b.rejected 中止全图。返回应激活的下游(拒绝=空)。 +func (o *Orchestrator) approvalNode(ctx context.Context, taskID string, n dsl.Node, b *board, tr *execTracer, outs []dsl.Edge) []string { + title := firstNonEmpty(cstr(n.Config, "title"), labelOf(n, "人工审批")) + prompt := firstNonEmpty(cstr(n.Config, "prompt"), "请审批是否继续执行后续步骤") + summary := prompt + if b.answer != "" { // 带上当前产出预览,便于审批人判断 + summary = prompt + "\n—— 当前产出预览 ——\n" + truncate(b.answer, 400) + } + + // 未接审批通道(单测/降级)→ 自动放行,避免无人应答卡死。 + if o.approval == nil { + tr.info("approval:"+n.ID, "approval", title, "未接审批通道,自动放行") + return targetsOf(outs) + } + + // 暂停:发待审事件(UI 据 kind=approval & phase=await 弹批准/拒绝)+ 置任务 waiting。 + tr.emit("approval:"+n.ID, "approval", "await", title, summary, 0) + o.setStatus(taskID, contract.TaskWaiting, title) + + dec, err := o.approval.WaitApproval(ctx, taskID, approvalTimeout) + switch { + case err != nil: // 超时 / ctx 取消 → 安全默认拒绝 + b.rejected = true + b.answer = "❌ 审批超时未决,已自动拒绝:" + title + tr.emit("approval:"+n.ID, "approval", "error", title, "审批超时,自动拒绝", 0) + return nil + case !dec.Approved: + b.rejected = true + note := firstNonEmpty(dec.Note, "审批人拒绝") + b.answer = "❌ 已被拒绝:" + note + tr.emit("approval:"+n.ID, "approval", "end", title, "拒绝:"+note, 0) + return nil + default: // 批准 → 恢复执行,放行下游 + o.setStatus(taskID, contract.TaskRunning, "审批通过,继续执行") + tr.emit("approval:"+n.ID, "approval", "end", title, "批准:"+firstNonEmpty(dec.Note, "放行"), 0) + return targetsOf(outs) + } +} + // targetsOf 取一组边的目标节点 ID(保持顺序)。 func targetsOf(edges []dsl.Edge) []string { out := make([]string, 0, len(edges)) diff --git a/sundynix-dispatcher/internal/eino/orchestrator.go b/sundynix-dispatcher/internal/eino/orchestrator.go index 11d8ed9..4019d4b 100644 --- a/sundynix-dispatcher/internal/eino/orchestrator.go +++ b/sundynix-dispatcher/internal/eino/orchestrator.go @@ -40,6 +40,15 @@ type StatusSink interface { PublishTaskStatus(taskID, status, detail string) error } +// ApprovalWaiter 阻塞等待审批节点的人工决定(由 NATS bus 实现;可为 nil → 审批节点自动放行)。 +type ApprovalWaiter interface { + WaitApproval(ctx context.Context, taskID string, timeout time.Duration) (*contract.ApprovalDecision, error) +} + +// errRejected 是审批节点拒绝(或超时)时图执行返回的哨兵错误:它是合法终态而非故障, +// Handle 据此判 rejected 并优雅收尾(不计熔断失败)。 +var errRejected = errors.New("approval rejected") + // LLM 是编排所需的语言模型能力(生产由 *llm.Pool 实现)。抽成接口便于测试注入假模型。 type LLM interface { Ready() bool @@ -56,17 +65,22 @@ type LLM interface { const toolCallTimeout = 3 * time.Second // taskExecTimeout 是单个任务整体执行上限;超时即判 timeout(状态机),避免无限期"运行中"。 -const taskExecTimeout = 3 * time.Minute +// 含 HITL 审批等待预算(approvalTimeout)+ 常规图执行;须 < bus 消费者 AckWait(15min) 以免重投。 +const taskExecTimeout = 10 * time.Minute + +// approvalTimeout 是单个审批节点等待人工决定的上限;超时安全默认拒绝(fail-safe)。 +const approvalTimeout = 5 * time.Minute // Orchestrator 把每个 DSL 任务动态编译为 Eino 图并执行(记忆召回 → 工具节点 → 注入 → 流式)。 type Orchestrator struct { - pool LLM - breaker *harness.CircuitBreaker - eval *harness.Evaluator - sink TokenSink - tools ToolCaller - exec ExecSink - status StatusSink // 任务生命周期状态回写(可为 nil) + pool LLM + breaker *harness.CircuitBreaker + eval *harness.Evaluator + sink TokenSink + tools ToolCaller + exec ExecSink + status StatusSink // 任务生命周期状态回写(可为 nil) + approval ApprovalWaiter // HITL 审批等待(可为 nil → 审批节点自动放行) turnMu sync.Mutex // 保护 turns(攒批计数,多任务 goroutine 共享) turns map[string]int // sessionID → 累计轮次,用于每 N 轮触发 consolidate @@ -74,9 +88,9 @@ type Orchestrator struct { // NewOrchestrator 持有依赖;图按任务的 DSL 在 Handle 内动态编译。 // exec 为执行可视化事件出口(可为 nil,则不发轨迹事件);eval 为自动化评测(可为 nil); -// status 为任务生命周期状态回写出口(可为 nil)。 -func NewOrchestrator(pool LLM, breaker *harness.CircuitBreaker, eval *harness.Evaluator, sink TokenSink, tools ToolCaller, exec ExecSink, status StatusSink) (*Orchestrator, error) { - return &Orchestrator{pool: pool, breaker: breaker, eval: eval, sink: sink, tools: tools, exec: exec, status: status}, nil +// status 为任务生命周期状态回写出口(可为 nil);approval 为 HITL 审批等待(可为 nil)。 +func NewOrchestrator(pool LLM, breaker *harness.CircuitBreaker, eval *harness.Evaluator, sink TokenSink, tools ToolCaller, exec ExecSink, status StatusSink, approval ApprovalWaiter) (*Orchestrator, error) { + return &Orchestrator{pool: pool, breaker: breaker, eval: eval, sink: sink, tools: tools, exec: exec, status: status, approval: approval}, nil } // setStatus 回写一次任务状态流转(status 为 nil 时静默跳过)。 @@ -144,6 +158,17 @@ func (o *Orchestrator) Handle(ctx context.Context, t *contract.Task) error { // 按 DSL 图执行:compose.Graph(EINO_COMPOSE=1)或自研 graph.go(默认);agent 节点流式回流 token。 answer, err := o.executeGraph(tctx, t, tr) + if errors.Is(err, errRejected) { + // HITL 拒绝:合法终态,非故障。收尾流 + 置 rejected,不计熔断、不重投。 + slog.InfoContext(ctx, "task rejected by approval", "task_id", t.ID) + if answer != "" { + _ = o.sink.PublishToken(t.ID, []byte(answer)) + } + _ = o.sink.CompleteStream(t.ID) + o.breaker.Report(true) // 拒绝是人为决策,不算后端失败 + o.setStatus(t.ID, contract.TaskRejected, truncate(answer, 120)) + return nil + } if err != nil { span.RecordError(err) span.SetStatus(codes.Error, err.Error()) diff --git a/sundynix-dispatcher/internal/nats/subscriber.go b/sundynix-dispatcher/internal/nats/subscriber.go index bebaecf..c9b4f13 100644 --- a/sundynix-dispatcher/internal/nats/subscriber.go +++ b/sundynix-dispatcher/internal/nats/subscriber.go @@ -81,6 +81,11 @@ func (s *Subscriber) PublishTaskStatus(taskID, status, detail string) error { }) } +// WaitApproval 让 Subscriber 满足 eino.ApprovalWaiter,阻塞等待审批节点的人工决定。 +func (s *Subscriber) WaitApproval(ctx context.Context, taskID string, timeout time.Duration) (*contract.ApprovalDecision, error) { + return s.inner.WaitApproval(ctx, taskID, timeout) +} + // RequestModelConfig 向控制面(Gateway)取当前激活的对话模型配置。 func (s *Subscriber) RequestModelConfig(ctx context.Context) (*contract.ModelConfig, error) { return s.inner.RequestConfig(ctx, contract.ConfigKindChat) diff --git a/sundynix-gateway/internal/handler/task_handler.go b/sundynix-gateway/internal/handler/task_handler.go index b4c6be8..dabdc2c 100644 --- a/sundynix-gateway/internal/handler/task_handler.go +++ b/sundynix-gateway/internal/handler/task_handler.go @@ -89,6 +89,34 @@ func (h *Handler) TaskStatus(c *gin.Context) { c.JSON(http.StatusOK, gin.H{"task_id": id, "status": status, "detail": detail}) } +// ApproveTask: POST /api/v1/tasks/:id/approve {approved, node?, note?} —— 人工审批决定(HITL)。 +// 把决定经 NATS 发给 dispatcher,解除审批节点的阻塞(批准放行 / 拒绝中止)。 +func (h *Handler) ApproveTask(c *gin.Context) { + id := c.Param("id") + var body struct { + Approved bool `json:"approved"` + Node string `json:"node"` + Note string `json:"note"` + } + if err := c.ShouldBindJSON(&body); err != nil { + c.JSON(http.StatusBadRequest, gin.H{"error": err.Error()}) + return + } + // 仅在任务确为等待审批时受理(幂等:重复/迟到的决定不报错,dispatcher 侧已只取首条)。 + if status, _ := h.db.GetTaskStatus(c.Request.Context(), id); status != contract.TaskWaiting { + c.JSON(http.StatusConflict, gin.H{"error": "任务当前非待审批状态", "status": status}) + return + } + if err := h.bus.PublishApproval(&contract.ApprovalDecision{ + TaskID: id, Node: body.Node, Approved: body.Approved, Note: body.Note, + By: userID(c), TS: time.Now().UnixMilli(), + }); err != nil { + c.JSON(http.StatusBadGateway, gin.H{"error": err.Error()}) + return + } + c.JSON(http.StatusOK, gin.H{"task_id": id, "approved": body.Approved}) +} + // StreamTask: 以 SSE 把 Token Stream 推给客户端。 // 优先从 Redis Stream 读(可回放 + 断点续传,根治连晚/重连丢 token);Redis 降级时回退 live NATS。 func (h *Handler) StreamTask(c *gin.Context) { diff --git a/sundynix-gateway/internal/nats/publisher.go b/sundynix-gateway/internal/nats/publisher.go index 181a8ea..4905d07 100644 --- a/sundynix-gateway/internal/nats/publisher.go +++ b/sundynix-gateway/internal/nats/publisher.go @@ -64,6 +64,11 @@ func (b *Bus) SubscribeTaskStatus(onEvent func(*contract.TaskStatusEvent)) (func return b.inner.SubscribeTaskStatus(onEvent) } +// PublishApproval 把一次人工审批决定发给 dispatcher(解除审批节点阻塞)。 +func (b *Bus) PublishApproval(dec *contract.ApprovalDecision) error { + return b.inner.PublishApproval(dec) +} + // ServeConfig 让网关作为配置控制面,响应某 kind 的配置请求。 func (b *Bus) ServeConfig(kind string, provide func() *contract.ModelConfig) (func() error, error) { return b.inner.ServeConfig(kind, provide) diff --git a/sundynix-gateway/internal/router/router.go b/sundynix-gateway/internal/router/router.go index 977d851..f455292 100644 --- a/sundynix-gateway/internal/router/router.go +++ b/sundynix-gateway/internal/router/router.go @@ -50,25 +50,26 @@ func New(db *store.Postgres, cache *store.Redis, bus *nats.Bus, blobStore *blob. // —— 受保护:owner 作用域业务,必须携带有效 JWT —— p := api.Group("", middleware.RequireAuth()) { - p.POST("/tasks", h.SubmitTask) // 解析 DSL 并 Publish 到 NATS(带已验证 uid) - p.GET("/tasks/:id", h.TaskStatus) // 任务生命周期状态(UI 轮询 submitted/running/done/failed/timeout) - p.PUT("/memory", h.SetMemory) // 偏好记忆登记(→ mcp-go memory_upsert) - p.GET("/memory", h.ListMemory) // 列出当前用户偏好(记忆面板) - p.DELETE("/memory", h.DeleteMemory) // 软删一条偏好(?key=) - p.GET("/kb/list", h.KbList) // 当前用户的知识库列表(owner 隔离) - p.POST("/kb/create", h.KbCreate) // 新建知识库 - p.POST("/kb/ingest", h.KbIngest) // 文本入库 - p.POST("/kb/ingest_file", h.KbIngestFile) // 文件入库 - p.POST("/kb/search", h.KbSearch) // 检索台 - p.GET("/kb/vault", h.KbVault) // 文库列表 - p.GET("/kb/doc", h.KbDoc) // 取单篇文档 - p.GET("/kb/links", h.KbLinks) // 某库双链 - p.POST("/kb/note", h.KbSaveNote) // 新建/编辑笔记 - p.GET("/kb/graph", h.KbGraph) // 知识图谱三元组 - p.GET("/agents", h.AgentList) // 我的编排列表(owner 隔离) - p.POST("/agents", h.AgentSave) // 保存/更新编排 - p.DELETE("/agents", h.AgentDelete) // 删除编排 - p.POST("/reports", h.GenerateReport) // 报告生成 + p.POST("/tasks", h.SubmitTask) // 解析 DSL 并 Publish 到 NATS(带已验证 uid) + p.GET("/tasks/:id", h.TaskStatus) // 任务生命周期状态(UI 轮询 submitted/running/done/failed/timeout/waiting/rejected) + p.POST("/tasks/:id/approve", h.ApproveTask) // HITL 人工审批决定(批准/拒绝) + p.PUT("/memory", h.SetMemory) // 偏好记忆登记(→ mcp-go memory_upsert) + p.GET("/memory", h.ListMemory) // 列出当前用户偏好(记忆面板) + p.DELETE("/memory", h.DeleteMemory) // 软删一条偏好(?key=) + p.GET("/kb/list", h.KbList) // 当前用户的知识库列表(owner 隔离) + p.POST("/kb/create", h.KbCreate) // 新建知识库 + p.POST("/kb/ingest", h.KbIngest) // 文本入库 + p.POST("/kb/ingest_file", h.KbIngestFile) // 文件入库 + p.POST("/kb/search", h.KbSearch) // 检索台 + p.GET("/kb/vault", h.KbVault) // 文库列表 + p.GET("/kb/doc", h.KbDoc) // 取单篇文档 + p.GET("/kb/links", h.KbLinks) // 某库双链 + p.POST("/kb/note", h.KbSaveNote) // 新建/编辑笔记 + p.GET("/kb/graph", h.KbGraph) // 知识图谱三元组 + p.GET("/agents", h.AgentList) // 我的编排列表(owner 隔离) + p.POST("/agents", h.AgentSave) // 保存/更新编排 + p.DELETE("/agents", h.AgentDelete) // 删除编排 + p.POST("/reports", h.GenerateReport) // 报告生成 p.GET("/billing", h.Billing) } diff --git a/sundynix-shared/bus/bus.go b/sundynix-shared/bus/bus.go index 939ff56..7f08bd9 100644 --- a/sundynix-shared/bus/bus.go +++ b/sundynix-shared/bus/bus.go @@ -314,6 +314,48 @@ func (b *Bus) SubscribeTaskStatus(onEvent func(*contract.TaskStatusEvent)) (unsu return sub.Unsubscribe, nil } +// ---- 人工审批(HITL,core NATS pub-sub)---- + +// PublishApproval 广播一次人工审批决定(网关在收到 UI 的批准/拒绝后调用)。 +func (b *Bus) PublishApproval(dec *contract.ApprovalDecision) error { + data, err := json.Marshal(dec) + if err != nil { + return err + } + return b.nc.Publish(contract.ApprovalSubject(dec.TaskID), data) +} + +// WaitApproval 阻塞等待某任务的人工审批决定,直到收到、ctx 取消或超时。 +// dispatcher 在审批节点调用:先订阅再等待(订阅早于决定到达,避免错过)。 +// 超时返回 (nil, error) —— 调用方据安全默认(拒绝)处理。 +func (b *Bus) WaitApproval(ctx context.Context, taskID string, timeout time.Duration) (*contract.ApprovalDecision, error) { + ch := make(chan *contract.ApprovalDecision, 1) + sub, err := b.nc.Subscribe(contract.ApprovalSubject(taskID), func(m *nats.Msg) { + var dec contract.ApprovalDecision + if json.Unmarshal(m.Data, &dec) == nil { + select { + case ch <- &dec: + default: // 已收到一条,丢弃后续重复决定 + } + } + }) + if err != nil { + return nil, fmt.Errorf("subscribe approval: %w", err) + } + defer func() { _ = sub.Unsubscribe() }() + + timer := time.NewTimer(timeout) + defer timer.Stop() + select { + case dec := <-ch: + return dec, nil + case <-timer.C: + return nil, fmt.Errorf("approval timeout after %s", timeout) + case <-ctx.Done(): + return nil, ctx.Err() + } +} + // ---- 配置控制面(core NATS request-reply + broadcast)---- // RequestConfig 向控制面(Gateway)请求某 kind 当前激活配置(chat/embedding)。 @@ -406,6 +448,9 @@ func (b *Bus) ConsumeTasks(ctx context.Context, h TaskHandler) (stop func(), err Durable: contract.ConsumerDurable, AckPolicy: jetstream.AckExplicitPolicy, FilterSubject: contract.SubjectTasksAll, + // HITL:审批节点会让 Handle 阻塞等人工决定(最长约 5 分钟), + // AckWait 必须覆盖「审批等待 + 图执行」总时长,否则消息在途未 ack 会被重投成重复任务。 + AckWait: 15 * time.Minute, }) if err != nil { return nil, fmt.Errorf("create consumer: %w", err) diff --git a/sundynix-shared/contract/task.go b/sundynix-shared/contract/task.go index fa20099..4e4178a 100644 --- a/sundynix-shared/contract/task.go +++ b/sundynix-shared/contract/task.go @@ -33,26 +33,46 @@ const ( // core NATS pub-sub(状态是幂等覆盖,丢一条由下一条纠正,无需持久化)。 // 注意:必须在 sundynix.tasks.> 之外,否则会被任务流捕获成"幽灵任务"自我放大。 SubjectTaskStatus = "sundynix.status.task" + + // 人工审批(HITL)决定回传前缀:实际 sundynix.approval.。 + // core NATS pub-sub:UI 点批准/拒绝 → 网关发到此 → dispatcher 解除审批节点的阻塞。 + SubjectApproval = "sundynix.approval" ) // 任务生命周期状态机:submitted(网关建任务)→ running(dispatcher 开跑) // → done / failed / timeout(dispatcher 收尾)。 +// HITL:执行到审批节点 → waiting(等人工决定)→ 批准回 running / 拒绝→rejected。 const ( TaskSubmitted = "submitted" TaskRunning = "running" TaskDone = "done" TaskFailed = "failed" TaskTimeout = "timeout" + TaskWaiting = "waiting" // 等待人工审批 + TaskRejected = "rejected" // 人工拒绝(或审批超时,安全默认拒绝) ) // TaskStatusEvent 是一次任务状态流转事件(经 SubjectTaskStatus 回流给网关)。 type TaskStatusEvent struct { TaskID string `json:"task_id"` - Status string `json:"status"` // running / done / failed / timeout - Detail string `json:"detail,omitempty"` // 失败原因等 + Status string `json:"status"` // running / done / failed / timeout / waiting / rejected + Detail string `json:"detail,omitempty"` // 失败原因 / 审批摘要等 TS int64 `json:"ts"` // unix 毫秒 } +// ApprovalSubject 返回某任务的人工审批决定回传主题。 +func ApprovalSubject(id string) string { return SubjectApproval + "." + id } + +// ApprovalDecision 是一次人工审批的结果(UI → 网关 → dispatcher)。 +type ApprovalDecision struct { + TaskID string `json:"task_id"` + Node string `json:"node,omitempty"` // 审批节点 id(可空:按 task 维度兜底匹配) + Approved bool `json:"approved"` // true=批准放行,false=拒绝中止 + Note string `json:"note,omitempty"` // 审批人备注 + By string `json:"by,omitempty"` // 审批人(用户 id) + TS int64 `json:"ts"` // unix 毫秒 +} + const ( // MetaUserID 是 Task.Meta 中承载已登录用户标识的键(用于偏好记忆召回)。 MetaUserID = "user_id"