From e42c77c8ae1e0dd779a75ca031cf84fdbbb64dfb Mon Sep 17 00:00:00 2001 From: mlogclub Date: Sun, 12 Apr 2026 13:36:01 +0800 Subject: [PATCH] feat: add AnalyzeConversation tool and related graph for analyzing conversation risks and summaries --- .../graphs/analyze_conversation_graph.go | 236 ++++++++++++++++++ .../graphs/analyze_conversation_graph_test.go | 53 ++++ .../internal/impl/factory/agent_factory.go | 21 +- internal/ai/runtime/registry/registry.go | 8 + internal/ai/runtime/service.go | 1 + .../tools/analyze_conversation_tool.go | 114 +++++++++ internal/pkg/toolx/builtin_tools.go | 4 + 7 files changed, 433 insertions(+), 4 deletions(-) create mode 100644 internal/ai/runtime/graphs/analyze_conversation_graph.go create mode 100644 internal/ai/runtime/graphs/analyze_conversation_graph_test.go create mode 100644 internal/ai/runtime/tools/analyze_conversation_tool.go diff --git a/internal/ai/runtime/graphs/analyze_conversation_graph.go b/internal/ai/runtime/graphs/analyze_conversation_graph.go new file mode 100644 index 0000000..fc70917 --- /dev/null +++ b/internal/ai/runtime/graphs/analyze_conversation_graph.go @@ -0,0 +1,236 @@ +package graphs + +import ( + "context" + "encoding/json" + "fmt" + "strings" + + "cs-agent/internal/models" + "cs-agent/internal/services" +) + +type AnalyzeConversationInput struct { + Goal string `json:"goal"` + ObservedIssue string `json:"observedIssue"` + NeedTicket bool `json:"needTicket"` + NeedHumanHandoff bool `json:"needHumanHandoff"` + NeedQualityCheck bool `json:"needQualityCheck"` + AdditionalContext string `json:"additionalContext"` +} + +type AnalyzeConversationResult struct { + Summary string `json:"summary"` + UserIntent string `json:"userIntent"` + RiskLevel string `json:"riskLevel"` + RiskSignals []string `json:"riskSignals,omitempty"` + RecommendedNextAction string `json:"recommendedNextAction"` + RecommendedQuestions []string `json:"recommendedQuestions,omitempty"` + ConversationFacts []string `json:"conversationFacts,omitempty"` +} + +type AnalyzeConversationGraph struct { + conversation *models.Conversation +} + +func NewAnalyzeConversationGraph(conversation *models.Conversation) *AnalyzeConversationGraph { + return &AnalyzeConversationGraph{conversation: conversation} +} + +func (g *AnalyzeConversationGraph) Run(_ context.Context, argumentsInJSON string) (string, error) { + if g == nil || g.conversation == nil { + return "", fmt.Errorf("analyze conversation graph not initialized") + } + input, err := g.parseInput(argumentsInJSON) + if err != nil { + return "", err + } + messages, _, _ := services.MessageService.FindByConversationIDCursor(g.conversation.ID, 0, 8, "", "") + result := buildAnalyzeConversationResult(g.conversation, messages, input) + buf, err := json.Marshal(result) + if err != nil { + return "", err + } + return string(buf), nil +} + +func (g *AnalyzeConversationGraph) parseInput(argumentsInJSON string) (AnalyzeConversationInput, error) { + var input AnalyzeConversationInput + if strings.TrimSpace(argumentsInJSON) == "" { + return input, nil + } + if err := json.Unmarshal([]byte(argumentsInJSON), &input); err != nil { + return input, fmt.Errorf("invalid analyze conversation arguments: %w", err) + } + input.Goal = strings.TrimSpace(input.Goal) + input.ObservedIssue = strings.TrimSpace(input.ObservedIssue) + input.AdditionalContext = strings.TrimSpace(input.AdditionalContext) + return input, nil +} + +func buildAnalyzeConversationResult(conversation *models.Conversation, messages []models.Message, input AnalyzeConversationInput) AnalyzeConversationResult { + joined := strings.ToLower(buildConversationCorpus(conversation, messages, input)) + signals := collectRiskSignals(joined, input) + intent := detectUserIntent(joined, input) + recommendedAction := recommendNextAction(intent, signals, input) + result := AnalyzeConversationResult{ + Summary: buildConversationSummary(conversation, messages, input), + UserIntent: intent, + RiskLevel: deriveRiskLevel(signals), + RiskSignals: signals, + RecommendedNextAction: recommendedAction, + RecommendedQuestions: recommendQuestions(intent, signals, input), + ConversationFacts: buildConversationAnalysisFacts(conversation, messages), + } + return result +} + +func buildConversationSummary(conversation *models.Conversation, messages []models.Message, input AnalyzeConversationInput) string { + parts := make([]string, 0, 4) + if conversation != nil && strings.TrimSpace(conversation.Subject) != "" { + parts = append(parts, "会话主题:"+strings.TrimSpace(conversation.Subject)) + } + if input.ObservedIssue != "" { + parts = append(parts, "当前问题:"+input.ObservedIssue) + } else if conversation != nil && strings.TrimSpace(conversation.LastMessageSummary) != "" { + parts = append(parts, "当前问题:"+strings.TrimSpace(conversation.LastMessageSummary)) + } + if digest := buildRecentMessageDigest(messages); digest != "" { + parts = append(parts, "最近对话:"+digest) + } + if input.AdditionalContext != "" { + parts = append(parts, "补充信息:"+input.AdditionalContext) + } + return strings.TrimSpace(strings.Join(parts, "\n")) +} + +func buildConversationCorpus(conversation *models.Conversation, messages []models.Message, input AnalyzeConversationInput) string { + parts := make([]string, 0, len(messages)+4) + if conversation != nil { + parts = append(parts, strings.TrimSpace(conversation.Subject)) + parts = append(parts, strings.TrimSpace(conversation.LastMessageSummary)) + } + parts = append(parts, input.Goal, input.ObservedIssue, input.AdditionalContext) + for i := range messages { + parts = append(parts, strings.TrimSpace(messages[i].Content)) + } + return strings.Join(parts, "\n") +} + +func collectRiskSignals(joined string, input AnalyzeConversationInput) []string { + signals := make([]string, 0, 6) + add := func(signal string) { + for _, item := range signals { + if item == signal { + return + } + } + signals = append(signals, signal) + } + if containsAny(joined, "投诉", "举报", "差评", "曝光", "媒体", "起诉", "律师", "12315") { + add("complaint_escalation") + } + if containsAny(joined, "退款", "赔偿", "损失", "扣款", "重复扣费", "金额") { + add("financial_risk") + } + if containsAny(joined, "生气", "愤怒", "垃圾", "太差", "一直没人", "再不处理") { + add("negative_sentiment") + } + if containsAny(joined, "人工", "转人工", "真人", "客服") || input.NeedHumanHandoff { + add("handoff_requested") + } + if containsAny(joined, "工单", "报障", "售后", "登记", "记录问题") || input.NeedTicket { + add("ticket_expected") + } + if input.NeedQualityCheck { + add("quality_review_requested") + } + return signals +} + +func detectUserIntent(joined string, input AnalyzeConversationInput) string { + switch { + case input.NeedHumanHandoff || containsAny(joined, "人工", "转人工", "真人"): + return "handoff_request" + case input.NeedTicket || containsAny(joined, "工单", "报障", "售后", "登记问题"): + return "ticket_request" + case containsAny(joined, "投诉", "举报", "差评", "赔偿"): + return "complaint" + default: + return "general_support" + } +} + +func deriveRiskLevel(signals []string) string { + if len(signals) == 0 { + return "low" + } + if containsSignal(signals, "complaint_escalation") || containsSignal(signals, "financial_risk") { + return "high" + } + if containsSignal(signals, "handoff_requested") || containsSignal(signals, "negative_sentiment") { + return "medium" + } + return "low" +} + +func recommendNextAction(intent string, signals []string, input AnalyzeConversationInput) string { + switch { + case input.NeedQualityCheck: + return "quality_review" + case containsSignal(signals, "handoff_requested") || intent == "handoff_request": + return "handoff_to_human" + case containsSignal(signals, "ticket_expected") || intent == "ticket_request": + return "prepare_ticket" + case containsSignal(signals, "complaint_escalation"): + return "handoff_to_human" + default: + return "continue_answering" + } +} + +func recommendQuestions(intent string, signals []string, input AnalyzeConversationInput) []string { + questions := make([]string, 0, 3) + if containsSignal(signals, "ticket_expected") && strings.TrimSpace(input.ObservedIssue) == "" { + questions = append(questions, "请进一步确认用户遇到的具体问题现象、报错信息和期望处理结果。") + } + if containsSignal(signals, "handoff_requested") { + questions = append(questions, "请确认用户是否明确要求人工客服,以及当前问题为何需要人工继续处理。") + } + if intent == "complaint" { + questions = append(questions, "请确认投诉点、影响范围和用户当前最希望解决的事项。") + } + return questions +} + +func buildConversationAnalysisFacts(conversation *models.Conversation, messages []models.Message) []string { + facts := make([]string, 0, 4) + if conversation != nil && strings.TrimSpace(conversation.Subject) != "" { + facts = append(facts, "会话主题:"+strings.TrimSpace(conversation.Subject)) + } + if conversation != nil && strings.TrimSpace(conversation.LastMessageSummary) != "" { + facts = append(facts, "最近摘要:"+strings.TrimSpace(conversation.LastMessageSummary)) + } + if digest := buildRecentMessageDigest(messages); digest != "" { + facts = append(facts, "最近消息:"+digest) + } + return facts +} + +func containsAny(value string, keywords ...string) bool { + for _, keyword := range keywords { + if keyword != "" && strings.Contains(value, strings.ToLower(keyword)) { + return true + } + } + return false +} + +func containsSignal(signals []string, target string) bool { + for _, item := range signals { + if item == target { + return true + } + } + return false +} diff --git a/internal/ai/runtime/graphs/analyze_conversation_graph_test.go b/internal/ai/runtime/graphs/analyze_conversation_graph_test.go new file mode 100644 index 0000000..7d5cbc7 --- /dev/null +++ b/internal/ai/runtime/graphs/analyze_conversation_graph_test.go @@ -0,0 +1,53 @@ +package graphs + +import ( + "testing" + + "cs-agent/internal/models" + "cs-agent/internal/pkg/enums" +) + +func TestBuildAnalyzeConversationResult_RecommendsHandoffForComplaint(t *testing.T) { + conversation := &models.Conversation{ + Subject: "用户投诉扣费异常", + LastMessageSummary: "用户反馈被重复扣费,并要求人工处理", + } + messages := []models.Message{ + {SenderType: enums.IMSenderTypeCustomer, Content: "你们重复扣费了,我要投诉并转人工"}, + } + + got := buildAnalyzeConversationResult(conversation, messages, AnalyzeConversationInput{ + NeedHumanHandoff: true, + }) + + if got.UserIntent != "handoff_request" { + t.Fatalf("expected handoff_request, got %q", got.UserIntent) + } + if got.RiskLevel != "high" { + t.Fatalf("expected high risk, got %q", got.RiskLevel) + } + if got.RecommendedNextAction != "handoff_to_human" { + t.Fatalf("expected handoff_to_human, got %q", got.RecommendedNextAction) + } +} + +func TestBuildAnalyzeConversationResult_RecommendsPrepareTicket(t *testing.T) { + conversation := &models.Conversation{ + Subject: "订单无法支付", + LastMessageSummary: "用户要求登记问题并尽快处理", + } + messages := []models.Message{ + {SenderType: enums.IMSenderTypeCustomer, Content: "麻烦帮我建个工单,订单一直支付失败"}, + } + + got := buildAnalyzeConversationResult(conversation, messages, AnalyzeConversationInput{ + NeedTicket: true, + }) + + if got.UserIntent != "ticket_request" { + t.Fatalf("expected ticket_request, got %q", got.UserIntent) + } + if got.RecommendedNextAction != "prepare_ticket" { + t.Fatalf("expected prepare_ticket, got %q", got.RecommendedNextAction) + } +} diff --git a/internal/ai/runtime/internal/impl/factory/agent_factory.go b/internal/ai/runtime/internal/impl/factory/agent_factory.go index 11f1505..5ba1ce9 100644 --- a/internal/ai/runtime/internal/impl/factory/agent_factory.go +++ b/internal/ai/runtime/internal/impl/factory/agent_factory.go @@ -109,16 +109,20 @@ func (f *AgentFactory) BuildCustomerServiceAgent(ctx context.Context, input Buil continue } serverCode, toolName := "", "" - if toolCode == toolx.BuiltinToolSearchToolCode { + switch toolCode { + case toolx.BuiltinToolSearchToolCode: serverCode = toolx.BuiltinToolCatalogServerCode toolName = toolx.BuiltinToolSearchToolName - } else if toolCode == toolx.GraphPrepareTicketDraftToolCode { + case toolx.GraphAnalyzeConversationToolCode: + serverCode = toolx.GraphToolCatalogServerCode + toolName = toolx.GraphAnalyzeConversationToolName + case toolx.GraphPrepareTicketDraftToolCode: serverCode = toolx.GraphToolCatalogServerCode toolName = toolx.GraphPrepareTicketDraftToolName - } else if toolCode == toolx.GraphCreateTicketConfirmToolCode { + case toolx.GraphCreateTicketConfirmToolCode: serverCode = toolx.GraphToolCatalogServerCode toolName = toolx.GraphCreateTicketConfirmToolName - } else if toolCode == toolx.GraphHandoffConversationToolCode { + case toolx.GraphHandoffConversationToolCode: serverCode = toolx.GraphToolCatalogServerCode toolName = toolx.GraphHandoffConversationToolName } @@ -222,6 +226,15 @@ func assembleAgentInstruction(aiAgent *models.AIAgent, selectedSkill *models.Ski 2. 如果工具返回 ready=false,优先根据 missingFields 和 followUpQuestions 继续追问,不要直接创建工单。 3. 如果工具返回 ready=true,再结合结果考虑调用 create_ticket_with_confirmation。 4. 该工具用于“整理草稿”,不代表已经创建工单。 +`)) + } + if hasToolCode(extraToolCodes, toolx.GraphAnalyzeConversationToolCode) { + appendixParts = append(appendixParts, strings.TrimSpace(` +当对话可能涉及投诉升级、退款赔偿、明显负面情绪、是否要建单、是否要转人工等复杂判断时,优先调用 analyze_conversation 这个 Graph Tool,并遵守以下规则: +1. 该工具用于输出结构化摘要、风险信号和下一步建议,不代表实际已经建单或转人工。 +2. 如果工具建议为 handoff_to_human,应先确认是否满足转人工条件,再考虑调用 handoff_to_human。 +3. 如果工具建议为 prepare_ticket,应优先调用 prepare_ticket_draft 或继续补充信息,而不是直接建单。 +4. 如果工具建议为 continue_answering,优先继续澄清和解答,不要过早升级动作。 `)) } if hasToolCode(extraToolCodes, toolx.GraphCreateTicketConfirmToolCode) { diff --git a/internal/ai/runtime/registry/registry.go b/internal/ai/runtime/registry/registry.go index 79bff08..3e13fda 100644 --- a/internal/ai/runtime/registry/registry.go +++ b/internal/ai/runtime/registry/registry.go @@ -59,6 +59,14 @@ func isAllowedToolCode(toolCode string, allowedToolCodes map[string]struct{}) bo if isAlwaysAllowedToolCode(toolCode) { return true } + if strings.TrimSpace(toolCode) == toolx.GraphAnalyzeConversationToolCode { + if _, ok := allowedToolCodes[toolx.GraphCreateTicketConfirmToolCode]; ok { + return true + } + if _, ok := allowedToolCodes[toolx.GraphHandoffConversationToolCode]; ok { + return true + } + } if strings.TrimSpace(toolCode) == toolx.GraphPrepareTicketDraftToolCode { _, ok := allowedToolCodes[toolx.GraphCreateTicketConfirmToolCode] return ok diff --git a/internal/ai/runtime/service.go b/internal/ai/runtime/service.go index 24182e6..738df19 100644 --- a/internal/ai/runtime/service.go +++ b/internal/ai/runtime/service.go @@ -19,6 +19,7 @@ func newService() *service { return &service{ runtime: engine.NewService(), registry: registry.NewRegistry( + tools.NewAnalyzeConversationTool(), tools.NewPrepareTicketDraftTool(), tools.NewCreateTicketGraphTool(), tools.NewHandoffGraphTool(), diff --git a/internal/ai/runtime/tools/analyze_conversation_tool.go b/internal/ai/runtime/tools/analyze_conversation_tool.go new file mode 100644 index 0000000..50e2293 --- /dev/null +++ b/internal/ai/runtime/tools/analyze_conversation_tool.go @@ -0,0 +1,114 @@ +package tools + +import ( + "context" + "fmt" + + "cs-agent/internal/ai/runtime/graphs" + "cs-agent/internal/ai/runtime/registry" + "cs-agent/internal/models" + "cs-agent/internal/pkg/toolx" + + einotool "github.com/cloudwego/eino/components/tool" + "github.com/cloudwego/eino/schema" + einojsonschema "github.com/eino-contrib/jsonschema" + orderedmap "github.com/wk8/go-ordered-map/v2" +) + +const ( + AnalyzeConversationToolCode = toolx.GraphAnalyzeConversationToolCode + AnalyzeConversationToolName = toolx.GraphAnalyzeConversationToolName +) + +type AnalyzeConversationTool struct { + conversation *models.Conversation +} + +func NewAnalyzeConversationTool() *AnalyzeConversationTool { + return &AnalyzeConversationTool{} +} + +func (t *AnalyzeConversationTool) Name() string { + return AnalyzeConversationToolName +} + +func (t *AnalyzeConversationTool) Code() string { + return AnalyzeConversationToolCode +} + +func (t *AnalyzeConversationTool) Enabled(ctx registry.Context) bool { + return ctx.Conversation != nil +} + +func (t *AnalyzeConversationTool) Build(ctx registry.Context) (einotool.BaseTool, error) { + if !t.Enabled(ctx) { + return nil, nil + } + return &AnalyzeConversationTool{conversation: ctx.Conversation}, nil +} + +func (t *AnalyzeConversationTool) Info(ctx context.Context) (*schema.ToolInfo, error) { + return &schema.ToolInfo{ + Name: AnalyzeConversationToolName, + Desc: "Graph Tool。用于整理当前对话摘要、识别投诉/资金/情绪等风险信号,并给出继续解答、建单或转人工的建议。", + ParamsOneOf: schema.NewParamsOneOfByJSONSchema(&einojsonschema.Schema{ + Version: einojsonschema.Version, + Type: "object", + Properties: orderedmap.New[string, *einojsonschema.Schema](orderedmap.WithInitialData( + orderedmap.Pair[string, *einojsonschema.Schema]{ + Key: "goal", + Value: &einojsonschema.Schema{ + Type: "string", + Description: "当前分析目标,例如判断是否需要转人工、是否需要建单、是否需要做风险质检。", + }, + }, + orderedmap.Pair[string, *einojsonschema.Schema]{ + Key: "observedIssue", + Value: &einojsonschema.Schema{ + Type: "string", + Description: "当前观察到的主要问题或诉求。", + }, + }, + orderedmap.Pair[string, *einojsonschema.Schema]{ + Key: "needTicket", + Value: &einojsonschema.Schema{ + Type: "boolean", + Description: "是否重点评估建单必要性。", + }, + }, + orderedmap.Pair[string, *einojsonschema.Schema]{ + Key: "needHumanHandoff", + Value: &einojsonschema.Schema{ + Type: "boolean", + Description: "是否重点评估转人工必要性。", + }, + }, + orderedmap.Pair[string, *einojsonschema.Schema]{ + Key: "needQualityCheck", + Value: &einojsonschema.Schema{ + Type: "boolean", + Description: "是否重点做风险/质检分析。", + }, + }, + orderedmap.Pair[string, *einojsonschema.Schema]{ + Key: "additionalContext", + Value: &einojsonschema.Schema{ + Type: "string", + Description: "补充上下文,例如你已发现的争议点、投诉点或业务限制。", + }, + }, + )), + }), + Extra: map[string]any{ + "toolCode": AnalyzeConversationToolCode, + "sourceType": "graph", + }, + }, nil +} + +func (t *AnalyzeConversationTool) InvokableRun(ctx context.Context, argumentsInJSON string, opts ...einotool.Option) (string, error) { + if t == nil || t.conversation == nil { + return "", fmt.Errorf("analyze conversation tool not initialized") + } + return graphs.NewAnalyzeConversationGraph(t.conversation).Run(ctx, argumentsInJSON) +} diff --git a/internal/pkg/toolx/builtin_tools.go b/internal/pkg/toolx/builtin_tools.go index bdd32ec..8c7fe44 100644 --- a/internal/pkg/toolx/builtin_tools.go +++ b/internal/pkg/toolx/builtin_tools.go @@ -11,6 +11,10 @@ const ( BuiltinSkillToolTitle = "加载专项技能说明" BuiltinSkillToolDescription = "用于加载当前命中的专项技能说明文档。仅在本轮已命中 Skill 时可用,适合将专项处理规则按需注入上下文。" GraphToolCatalogServerCode = "graph" + GraphAnalyzeConversationToolCode = "graph/analyze_conversation" + GraphAnalyzeConversationToolName = "analyze_conversation" + GraphAnalyzeConversationToolTitle = "分析对话风险与摘要" + GraphAnalyzeConversationToolDescription = "Graph Tool。用于整理当前对话摘要、识别风险信号,并给出继续解答、建单或转人工的建议。" GraphPrepareTicketDraftToolCode = "graph/prepare_ticket_draft" GraphPrepareTicketDraftToolName = "prepare_ticket_draft" GraphPrepareTicketDraftToolTitle = "整理工单草稿"