/** Group 调度策略:为当前消息上下文生成一组有依赖关系的成员投递节点。 */ import type { GroupMessage } from "@/types/group/Group.js"; import type { Agent } from "@/agent/Agent.js"; import type { AgentModel } from "@/agent/AgentModel.js"; import type { ModelClient, ModelContent, ModelJsonValue, ModelMessage, } from "@downcity/type"; import { generate_model } from "@executor/model/ModelGenerate.js"; import { z } from "zod"; /** 一个投递节点内成员的发言方式。 */ export type DispatchResponseMode = "single" | "parallel"; /** 本次调度的触发来源。 */ export type DispatchTrigger = "user" | "auto"; /** 调度图中的一个成员投递节点。 */ export interface DispatchNode { /** 当前调度图内的稳定节点标识。 */ readonly node_id: string; /** 当前节点允许接收消息的成员标识。 */ readonly member_ids: readonly string[]; /** 当前节点只运行首个成员,还是允许全部成员并行运行。 */ readonly response_mode: DispatchResponseMode; /** 当前节点依赖的前置节点标识;依赖完成后才能投递。 */ readonly depends_on_node_ids: readonly string[]; /** 注入成员 AgentSession 的发言约束。 */ readonly instruction: string; } /** 一次调度生成的当前最佳响应图。 */ export interface DispatchDecision { /** 当前调度图中的投递节点。 */ readonly nodes: readonly DispatchNode[]; /** 当前图完成后是否允许唯一的 GroupSession auto dispatch 继续判断。 */ readonly terminal: boolean; } /** Group 调度策略协议。 */ export interface DispatchStrategy { /** 根据触发来源、当前消息、历史和成员关系生成当前最佳响应图。 */ decide_dispatch(input: { /** 本次调度由用户消息触发,还是由 GroupSession 自动收口触发。 */ readonly trigger: DispatchTrigger; /** 当前待调度的 Group 消息;auto dispatch 时为本批新增消息中的最后一条。 */ readonly message: GroupMessage; /** 当前 auto dispatch 尚未消费的新增消息。 */ readonly pending_messages: readonly GroupMessage[]; /** 当前 Group 已存在的消息历史。 */ readonly messages: readonly GroupMessage[]; /** 当前 Group 成员快照。 */ readonly members: readonly Agent[]; /** 当前 Dispatch Session Turn 的取消信号。 */ readonly abort_signal: AbortSignal; }): Promise | DispatchDecision; } /** AI 调度策略的构造参数。 */ export interface AiDispatchStrategyOptions { /** Group 用于理解群聊意图并生成调度图的模型。 */ readonly model?: AgentModel; } /** AI 调度工具的最小输入协议;外层阶段串行,阶段内成员并行。 */ const dispatch_group_input_schema = z.object({ steps: z.array(z.array(z.string().trim().min(1)).min(1).max(16)).max(16), /** 当前响应图完成后是否交给唯一的 auto dispatch 继续判断。 */ next: z.enum(["stop", "continue"]), }); const max_dispatch_model_steps = 3; /** 使用 Group.model 理解群聊意图并通过内部 tool call 生成成员投递决定。 */ export class AiDispatchStrategy implements DispatchStrategy { private readonly model?: ModelClient; constructor(options: AiDispatchStrategyOptions = {}) { this.model = options.model; } async decide_dispatch(input: { readonly trigger: DispatchTrigger; readonly message: GroupMessage; readonly pending_messages: readonly GroupMessage[]; readonly messages: readonly GroupMessage[]; readonly members: readonly Agent[]; readonly abort_signal: AbortSignal; }): Promise { if (!this.model) { throw new Error("Group requires a configured model for dispatch"); } try { const messages: ModelMessage[] = [ { role: "system", content: [{ type: "text", text: "你是群聊消息调度器。你只能调用 dispatch_group,不回答用户问题。普通任务默认只选择一个最合适的成员;只有用户明确要求多人分别回答,或确实存在并行的独立工作时,才把多个成员放在同一阶段。外层 steps 阶段按顺序执行,后一个阶段等待前一个阶段完成。只有真实存在的成员才能被选择。" }] }, { role: "user", content: [{ type: "text", text: build_dispatch_prompt(input) }] }, ]; let protocol_error = "Group dispatch model did not call dispatch_group"; for (let step_index = 0; step_index < max_dispatch_model_steps; step_index += 1) { const result = await generate_model(this.model, { messages, tools: [{ name: "dispatch_group", description: "提交当前 Group 的成员投递路径。只调用一次;不要输出普通文本。没有成员需要回复时使用空 steps 和 next=stop。", input_schema: z.toJSONSchema(dispatch_group_input_schema) as Record, }], tool_choice: { type: "tool", tool_name: "dispatch_group" }, max_output_tokens: 600, reasoning: { enabled: false }, }, input.abort_signal); const dispatch_call = result.tool_calls.length === 1 && result.tool_calls[0]?.tool_name === "dispatch_group" ? result.tool_calls[0] : undefined; if (dispatch_call) { try { return normalize_dispatch_tool_input(dispatch_call.input, input.members); } catch (error) { protocol_error = error instanceof Error ? error.message : String(error); } } else { protocol_error = "Group dispatch model did not call dispatch_group exactly once"; } if (step_index + 1 >= max_dispatch_model_steps) break; append_dispatch_retry_messages(messages, result.text, result.tool_calls, protocol_error); } throw new Error(`${protocol_error} after ${max_dispatch_model_steps} attempts`); } catch (error) { const detail = error instanceof Error ? error.message : String(error); throw new Error(`Group dispatch model failed: ${detail}`, { cause: error }); } } } /** 把协议错误作为下一模型 step 的显式上下文,允许模型在同一调度 Turn 内纠正。 */ function append_dispatch_retry_messages( messages: ModelMessage[], text: string, tool_calls: readonly { type: "tool_call"; tool_call_id: string; tool_name: string; input: ModelJsonValue }[], error: string, ): void { const assistant_content: ModelContent[] = [ ...(text.trim() ? [{ type: "text" as const, text }] : []), ...tool_calls, ]; if (assistant_content.length > 0) { messages.push({ role: "assistant", content: assistant_content }); } if (tool_calls.length > 0) { messages.push({ role: "tool", content: tool_calls.map((tool_call) => ({ type: "tool_result" as const, tool_call_id: tool_call.tool_call_id, tool_name: tool_call.tool_name, outcome: "failed" as const, content: [{ type: "text" as const, text: error }], })), }); } messages.push({ role: "user", content: [{ type: "text", text: `上一响应不符合调度协议:${error}。请只调用一次 dispatch_group。`, }], }); } /** 为 AI 调度器构造稳定、有限长度的群聊上下文。 */ function build_dispatch_prompt(input: { readonly trigger: DispatchTrigger; readonly message: GroupMessage; readonly pending_messages: readonly GroupMessage[]; readonly messages: readonly GroupMessage[]; readonly members: readonly Agent[]; }): string { const members = input.members.map((member) => `- ${member.id}`).join("\n"); const history = input.messages.slice(-20).map((message) => `${message.sender_type}:${message.sender_id}: ${message.text}`).join("\n"); return [ "成员:", members, "最近群聊:", history || "(无)", "本次新增消息:", input.pending_messages.map((message) => `${message.sender_type}:${message.sender_id}: ${message.text}`).join("\n") || "(无)", `调度触发来源:${input.trigger}`, "当前消息:", `${input.message.sender_type}:${input.message.sender_id}: ${input.message.text}`, "请调用 dispatch_group,输入 steps 和 next。steps 是二维数组:外层阶段按顺序执行,同一阶段数组内的成员并行执行;没有成员需要回复时使用空 steps 和 next=stop。不要输出普通文本。", ].join("\n"); } /** 将 AI tool call 编译为内部执行图,并确保只引用真实成员。 */ function normalize_dispatch_tool_input(value: unknown, members: readonly Agent[]): DispatchDecision { const parsed = dispatch_group_input_schema.parse(value); const valid_ids = new Set(members.map((member) => member.id)); const nodes = parsed.steps.map((member_ids, index) => { const unique_member_ids = [...new Set(member_ids)]; if (unique_member_ids.length !== member_ids.length) { throw new Error(`Group dispatch selected member more than once in step ${index}`); } for (const member_id of unique_member_ids) { if (!valid_ids.has(member_id)) { throw new Error(`Group dispatch selected unknown member: ${member_id}`); } } return { node_id: `step-${index}`, member_ids: unique_member_ids, response_mode: unique_member_ids.length > 1 ? "parallel" as const : "single" as const, depends_on_node_ids: index === 0 ? [] : [`step-${index - 1}`], instruction: "你已被 Group 调度选中,请直接针对当前消息给出回复,并且只代表自己发言。", }; }); if (nodes.length === 0 && parsed.next === "continue") { throw new Error("Group dispatch cannot continue without any member step"); } return { nodes, terminal: parsed.next === "stop", }; }