后端: 1. 日程暂存接口——前端拖拽调整后保存到 Redis 快照 - api/agent.go:新增 SaveScheduleState handler,解析绝对时间格式请求体,3 秒超时保护 - routers/routers.go:注册 POST /schedule-state - model/agent.go:新增 SaveScheduleStatePlacedItem / SaveScheduleStateRequest 结构体 - respond/respond.go:新增 5 个排程状态错误码(40058~40062) - 新增 service/agentsvc/agent_schedule_state.go:Load 快照 → ApplyPlacedItems → Save 回 Redis,校验归属 - 新增 newAgent/conv/schedule_state_apply.go:ApplyPlacedItems 绝对坐标→相对 day_index 转换,去重/坐标/嵌入关系校验 2. SchedulePersistor 持久化层全面下线 - 删除 newAgent/conv/schedule_persist.go(280 行,DiffScheduleState → applyChange → 事务写库整条链路) - model/state_store.go:移除 SchedulePersistor 接口 - model/graph_run_state.go / node/execute.go / node/agent_nodes.go / service/agent.go / service/agent_newagent.go / cmd/start.go:移除 SchedulePersistor 字段、参数、注入六处 3. schedule_completed 事件推送——deliver 节点排程完毕信号 - model/common_state.go:新增 HasScheduleChanges 标记,ResetForNextRun 清理 - node/execute.go / node/rough_build.go:写工具和粗排成功后置 HasScheduleChanges=true - node/deliver.go:IsCompleted && HasScheduleChanges 时调用 EmitScheduleCompleted - stream/emitter.go:新增 EmitScheduleCompleted 方法 - stream/openai.go:新增 StreamExtraKindScheduleCompleted + NewScheduleCompletedExtra 4. 预览接口补全 task_class_id - model/agent.go:GetSchedulePlanPreviewResponse 新增 TaskClassIDs - model/schedule.go:HybridScheduleEntry 新增 TaskClassID - conv/schedule_preview.go / service/agent_schedule_preview.go / service/schedule.go:三处透传填充 前端: 5. 排程完毕卡片 + 精排弹窗集成 - 新增 api/schedule_agent.ts:getSchedulePreview / saveScheduleState / applyBatchIntoSchedule - types/dashboard.ts:新增 HybridScheduleEntry / SchedulePreviewData / PlacedItem 类型 - components/dashboard/AssistantPanel.vue:监听 schedule_completed 事件异步拉取排程渲染卡片,集成 ScheduleResultCard + ScheduleFineTuneModal;confirm 交互从文本消息改为 resume 协议(approve/reject/cancel) 6. ToolTracePrototypeView 原型页新增日程小卡片 + 拖拽编排弹窗演示 7. DashboardView import 区域尺寸微调
249 lines
8.0 KiB
Go
249 lines
8.0 KiB
Go
package newagentnode
|
||
|
||
import (
|
||
"context"
|
||
"fmt"
|
||
"strings"
|
||
"time"
|
||
|
||
"github.com/cloudwego/eino/schema"
|
||
|
||
infrallm "github.com/LoveLosita/smartflow/backend/infra/llm"
|
||
newagentmodel "github.com/LoveLosita/smartflow/backend/newAgent/model"
|
||
newagentprompt "github.com/LoveLosita/smartflow/backend/newAgent/prompt"
|
||
newagentstream "github.com/LoveLosita/smartflow/backend/newAgent/stream"
|
||
)
|
||
|
||
const (
|
||
deliverStageName = "deliver"
|
||
deliverStatusBlockID = "deliver.status"
|
||
deliverSpeakBlockID = "deliver.speak"
|
||
)
|
||
|
||
// DeliverNodeInput 描述交付节点单轮运行所需的最小依赖。
|
||
//
|
||
// 职责边界:
|
||
// 1. 只负责生成交付总结并推送给用户,不负责后续流程推进;
|
||
// 2. RuntimeState 提供计划步骤和执行状态;
|
||
// 3. ConversationContext 提供执行阶段的对话历史;
|
||
// 4. 交付完成后标记流程结束。
|
||
type DeliverNodeInput struct {
|
||
RuntimeState *newagentmodel.AgentRuntimeState
|
||
ConversationContext *newagentmodel.ConversationContext
|
||
Client *infrallm.Client
|
||
ChunkEmitter *newagentstream.ChunkEmitter
|
||
ThinkingEnabled bool // 是否开启 thinking,由 config.yaml 的 agent.thinking.deliver 注入
|
||
CompactionStore newagentmodel.CompactionStore // 上下文压缩持久化
|
||
PersistVisibleMessage newagentmodel.PersistVisibleMessageFunc
|
||
}
|
||
|
||
// RunDeliverNode 执行一轮交付节点逻辑。
|
||
//
|
||
// 核心职责:
|
||
// 1. 调 LLM 基于原始计划 + 执行历史生成交付总结;
|
||
// 2. 伪流式推送总结给用户;
|
||
// 3. 写入对话历史,保证上下文连续;
|
||
// 4. 标记流程结束。
|
||
//
|
||
// 降级策略:
|
||
// 1. LLM 调用失败时,回退到机械格式化总结,不中断流程;
|
||
// 2. 机械总结包含计划步骤列表和完成进度。
|
||
func RunDeliverNode(ctx context.Context, input DeliverNodeInput) error {
|
||
runtimeState, conversationContext, emitter, err := prepareDeliverNodeInput(input)
|
||
if err != nil {
|
||
return err
|
||
}
|
||
flowState := runtimeState.EnsureCommonState()
|
||
|
||
// 1. 推送交付阶段状态,让前端知道正在生成总结。
|
||
if err := emitter.EmitStatus(
|
||
deliverStatusBlockID,
|
||
deliverStageName,
|
||
"summarizing",
|
||
"正在生成交付总结。",
|
||
false,
|
||
); err != nil {
|
||
return fmt.Errorf("交付阶段状态推送失败: %w", err)
|
||
}
|
||
|
||
// 2. 调 LLM 生成交付总结。
|
||
summary := generateDeliverSummary(ctx, input.Client, flowState, conversationContext, input.ThinkingEnabled, input.CompactionStore, emitter)
|
||
|
||
// 2.1 排程完毕卡片信号:
|
||
// 1. 仅在流程正常完成且确实产生过日程变更(粗排或写工具)时推送;
|
||
// 2. 前端收到 kind=schedule_completed 后,自行用对话 ID 调用现有接口拉取排程数据渲染卡片;
|
||
// 3. 不携带 Redis key 或排程数据,保持信号职责单一。
|
||
if flowState.IsCompleted() && flowState.HasScheduleChanges {
|
||
_ = emitter.EmitScheduleCompleted(deliverStatusBlockID, deliverStageName)
|
||
}
|
||
|
||
// 3. 伪流式推送总结。
|
||
if strings.TrimSpace(summary) != "" {
|
||
msg := schema.AssistantMessage(summary, nil)
|
||
if err := emitter.EmitPseudoAssistantText(
|
||
ctx,
|
||
deliverSpeakBlockID,
|
||
deliverStageName,
|
||
summary,
|
||
newagentstream.DefaultPseudoStreamOptions(),
|
||
); err != nil {
|
||
return fmt.Errorf("交付总结推送失败: %w", err)
|
||
}
|
||
conversationContext.AppendHistory(msg)
|
||
persistVisibleAssistantMessage(ctx, input.PersistVisibleMessage, flowState, msg)
|
||
}
|
||
|
||
// 4. 推送最终完成状态。
|
||
_ = emitter.EmitStatus(
|
||
deliverStatusBlockID,
|
||
deliverStageName,
|
||
"done",
|
||
"本轮流程已结束。",
|
||
true,
|
||
)
|
||
|
||
return nil
|
||
}
|
||
|
||
// generateDeliverSummary 尝试调用 LLM 生成交付总结,失败时降级到机械格式化。
|
||
func generateDeliverSummary(
|
||
ctx context.Context,
|
||
client *infrallm.Client,
|
||
flowState *newagentmodel.CommonState,
|
||
conversationContext *newagentmodel.ConversationContext,
|
||
thinkingEnabled bool,
|
||
compactionStore newagentmodel.CompactionStore,
|
||
emitter *newagentstream.ChunkEmitter,
|
||
) string {
|
||
if flowState != nil {
|
||
switch {
|
||
case flowState.IsAborted():
|
||
return normalizeSpeak(buildAbortSummary(flowState))
|
||
case flowState.IsExhaustedTerminal():
|
||
return normalizeSpeak(buildExhaustedSummary(flowState))
|
||
}
|
||
}
|
||
|
||
if client == nil {
|
||
return buildMechanicalSummary(flowState)
|
||
}
|
||
|
||
messages := newagentprompt.BuildDeliverMessages(flowState, conversationContext)
|
||
messages = compactUnifiedMessagesIfNeeded(ctx, messages, UnifiedCompactInput{
|
||
Client: client,
|
||
CompactionStore: compactionStore,
|
||
FlowState: flowState,
|
||
Emitter: emitter,
|
||
StageName: deliverStageName,
|
||
StatusBlockID: deliverStatusBlockID,
|
||
})
|
||
logNodeLLMContext(deliverStageName, "summarizing", flowState, messages)
|
||
result, err := client.GenerateText(
|
||
ctx,
|
||
messages,
|
||
infrallm.GenerateOptions{
|
||
Temperature: 0.5,
|
||
MaxTokens: 800,
|
||
Thinking: resolveThinkingMode(thinkingEnabled),
|
||
Metadata: map[string]any{
|
||
"stage": deliverStageName,
|
||
},
|
||
},
|
||
)
|
||
if err != nil || result == nil || strings.TrimSpace(result.Text) == "" {
|
||
return buildMechanicalSummary(flowState)
|
||
}
|
||
|
||
return normalizeSpeak(result.Text)
|
||
}
|
||
|
||
// buildAbortSummary 生成“流程已终止”的统一交付文案。
|
||
//
|
||
// 说明:
|
||
// 1. 第二轮开始,abort 的用户可见文案由终止方提前写入 CommonState;
|
||
// 2. deliver 不再重新猜测或改写业务异常,只做最终收口;
|
||
// 3. 若历史快照缺失 user_message,则回退到一份通用说明,避免前端收到空白结果。
|
||
func buildAbortSummary(state *newagentmodel.CommonState) string {
|
||
if state == nil || state.TerminalOutcome == nil {
|
||
return "本轮流程已终止。"
|
||
}
|
||
if msg := strings.TrimSpace(state.TerminalOutcome.UserMessage); msg != "" {
|
||
return msg
|
||
}
|
||
return "本轮流程已终止,请根据当前提示检查后再继续。"
|
||
}
|
||
|
||
// buildExhaustedSummary 生成“轮次耗尽”的统一收口文案。
|
||
func buildExhaustedSummary(state *newagentmodel.CommonState) string {
|
||
if state == nil {
|
||
return "本轮执行已达到安全轮次上限,当前先停止继续操作。"
|
||
}
|
||
|
||
prefix := "本轮执行已达到安全轮次上限,当前先停止继续操作。"
|
||
if state.TerminalOutcome != nil && strings.TrimSpace(state.TerminalOutcome.UserMessage) != "" {
|
||
prefix = strings.TrimSpace(state.TerminalOutcome.UserMessage)
|
||
}
|
||
if !state.HasPlan() {
|
||
return prefix
|
||
}
|
||
return prefix + "\n\n" + strings.TrimSpace(buildMechanicalSummary(state))
|
||
}
|
||
|
||
// buildMechanicalSummary 在 LLM 不可用时,机械拼接一份最小可用总结。
|
||
func buildMechanicalSummary(state *newagentmodel.CommonState) string {
|
||
if state == nil {
|
||
return "任务流程已结束。"
|
||
}
|
||
|
||
var sb strings.Builder
|
||
current, total := state.PlanProgress()
|
||
|
||
if !state.HasPlan() {
|
||
return "任务流程已结束。"
|
||
}
|
||
|
||
if state.IsExhaustedTerminal() {
|
||
sb.WriteString(fmt.Sprintf("任务因执行轮次耗尽提前结束,已完成 %d/%d 步。\n", current, total))
|
||
} else {
|
||
sb.WriteString("所有计划步骤已执行完毕。\n")
|
||
}
|
||
|
||
sb.WriteString("\n执行情况:\n")
|
||
for i, step := range state.PlanSteps {
|
||
marker := "[ ]"
|
||
if i < current {
|
||
marker = "[x]"
|
||
}
|
||
sb.WriteString(fmt.Sprintf("%s %s\n", marker, strings.TrimSpace(step.Content)))
|
||
}
|
||
|
||
if state.IsExhaustedTerminal() && current < total {
|
||
sb.WriteString("\n如需继续完成剩余步骤,可以告诉我继续。")
|
||
}
|
||
|
||
return sb.String()
|
||
}
|
||
|
||
// prepareDeliverNodeInput 校验并准备交付节点的运行态依赖。
|
||
func prepareDeliverNodeInput(input DeliverNodeInput) (
|
||
*newagentmodel.AgentRuntimeState,
|
||
*newagentmodel.ConversationContext,
|
||
*newagentstream.ChunkEmitter,
|
||
error,
|
||
) {
|
||
if input.RuntimeState == nil {
|
||
return nil, nil, nil, fmt.Errorf("deliver node: runtime state 不能为空")
|
||
}
|
||
|
||
input.RuntimeState.EnsureCommonState()
|
||
if input.ConversationContext == nil {
|
||
input.ConversationContext = newagentmodel.NewConversationContext("")
|
||
}
|
||
if input.ChunkEmitter == nil {
|
||
input.ChunkEmitter = newagentstream.NewChunkEmitter(
|
||
newagentstream.NoopPayloadEmitter(), "", "", time.Now().Unix(),
|
||
)
|
||
}
|
||
return input.RuntimeState, input.ConversationContext, input.ChunkEmitter, nil
|
||
}
|