18c9354095
- 注入数据库、运行时配置、统一响应、文件存储和平台 AI 能力,补充业务读写工具与客户快捷操作契约。 - 移除模块内重复的组织、客户、工单、标签、技能、旧工作流、MCP 和迁移实现,将身份权限与业务主体交由宿主管理。 - 使用 libSQL 重构向量存储,并完善图片消息、访客身份、排队调度、企业微信和支持聊天页面。 - 统一 HTTP、DTO 与 WebSocket 的 snake_case 协议,补齐模块初始化、业务动作和公共载荷等回归测试。
77 lines
2.5 KiB
Go
77 lines
2.5 KiB
Go
package runtime
|
|
|
|
import (
|
|
"context"
|
|
"fmt"
|
|
"strings"
|
|
|
|
applicationruntime "code.tczkiot.com/wlw/ai-agent/internal/ai/application/runtime"
|
|
"code.tczkiot.com/wlw/ai-agent/internal/ai/runtime/graphs"
|
|
"code.tczkiot.com/wlw/ai-agent/internal/models"
|
|
)
|
|
|
|
type runtimeReplyExecutor struct{}
|
|
|
|
type runtimeReplyRunInput struct {
|
|
Conversation models.Conversation
|
|
Message models.Message
|
|
AIAgent models.AIAgent
|
|
}
|
|
|
|
type runtimeReplyResumeInput struct {
|
|
Conversation models.Conversation
|
|
Message models.Message
|
|
AIAgent models.AIAgent
|
|
PendingInterrupt *models.ConversationInterrupt
|
|
}
|
|
|
|
func newRuntimeReplyExecutor() *runtimeReplyExecutor {
|
|
return &runtimeReplyExecutor{}
|
|
}
|
|
|
|
func (e *runtimeReplyExecutor) Run(ctx context.Context, input runtimeReplyRunInput) (*applicationruntime.RunResult, error) {
|
|
config, err := applicationruntime.ResolveRuntimeAIConfig(ctx, input.AIAgent.AIConfigID)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
// The trigger layer may enrich an anonymous channel conversation with a
|
|
// business subject resolved from the current message or recent history.
|
|
// Run the already validated objects so that card/device identity is not lost
|
|
// by reloading the original guest ownership record from the database.
|
|
summary, err := applicationruntime.DefaultAgentApplicationService.RunPrepared(ctx, applicationruntime.RunInput{
|
|
Conversation: input.Conversation,
|
|
UserMessage: input.Message,
|
|
AIAgent: input.AIAgent,
|
|
AIConfig: *config,
|
|
})
|
|
return summary, err
|
|
}
|
|
|
|
func (e *runtimeReplyExecutor) ResumePendingInterrupt(ctx context.Context, input runtimeReplyResumeInput) (*applicationruntime.RunResult, error) {
|
|
if input.PendingInterrupt == nil {
|
|
return nil, fmt.Errorf("pending interrupt is required")
|
|
}
|
|
config, err := applicationruntime.ResolveRuntimeAIConfig(ctx, input.AIAgent.AIConfigID)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
summary, err := applicationruntime.DefaultAgentApplicationService.ResumePrepared(ctx, applicationruntime.ResumeInput{
|
|
Conversation: input.Conversation,
|
|
UserMessage: input.Message,
|
|
AIAgent: input.AIAgent,
|
|
AIConfig: *config,
|
|
CheckPointID: strings.TrimSpace(input.PendingInterrupt.CheckPointID),
|
|
ResumeData: map[string]string{
|
|
strings.TrimSpace(input.PendingInterrupt.InterruptID): strings.TrimSpace(input.Message.Content),
|
|
},
|
|
})
|
|
return summary, err
|
|
}
|
|
|
|
func expiredInterruptSummary() *applicationruntime.RunResult {
|
|
return &applicationruntime.RunResult{
|
|
Status: "expired",
|
|
ReplyText: graphs.ConfirmationExpiredReply,
|
|
}
|
|
}
|