Files
ai-agent/internal/repositories/conversation_interrupt_repository.go
T
t 18c9354095 refactor: 将客服后端重构为宿主可嵌入模块
- 注入数据库、运行时配置、统一响应、文件存储和平台 AI 能力,补充业务读写工具与客户快捷操作契约。

- 移除模块内重复的组织、客户、工单、标签、技能、旧工作流、MCP 和迁移实现,将身份权限与业务主体交由宿主管理。

- 使用 libSQL 重构向量存储,并完善图片消息、访客身份、排队调度、企业微信和支持聊天页面。

- 统一 HTTP、DTO 与 WebSocket 的 snake_case 协议,补齐模块初始化、业务动作和公共载荷等回归测试。
2026-08-28 22:23:13 +08:00

102 lines
3.4 KiB
Go

package repositories
import (
"code.tczkiot.com/wlw/ai-agent/internal/models"
"github.com/mlogclub/simple/sqls"
"gorm.io/gorm"
)
var ConversationInterruptRepository = newConversationInterruptRepository()
func newConversationInterruptRepository() *conversationInterruptRepository {
return &conversationInterruptRepository{}
}
type conversationInterruptRepository struct{}
func (r *conversationInterruptRepository) Get(db *gorm.DB, id int64) *models.ConversationInterrupt {
ret := &models.ConversationInterrupt{}
if err := db.First(ret, "id = ?", id).Error; err != nil {
return nil
}
return ret
}
func (r *conversationInterruptRepository) GetByCheckPointID(db *gorm.DB, checkPointID string) *models.ConversationInterrupt {
ret := &models.ConversationInterrupt{}
if err := db.Where("check_point_id = ?", checkPointID).First(ret).Error; err != nil {
return nil
}
return ret
}
func (r *conversationInterruptRepository) FindLatestPendingByConversationID(db *gorm.DB, conversationID int64) *models.ConversationInterrupt {
ret := &models.ConversationInterrupt{}
if err := db.Where("conversation_id = ? AND status = ?", conversationID, "pending").Order("id DESC").First(ret).Error; err != nil {
return nil
}
return ret
}
func (r *conversationInterruptRepository) FindByAgentRunIDs(db *gorm.DB, agentRunIDs []int64) []models.ConversationInterrupt {
if len(agentRunIDs) == 0 {
return []models.ConversationInterrupt{}
}
var items []models.ConversationInterrupt
if err := db.Where("agent_run_id IN ?", agentRunIDs).Find(&items).Error; err != nil {
return []models.ConversationInterrupt{}
}
return items
}
func (r *conversationInterruptRepository) Find(db *gorm.DB, cnd *sqls.Cnd) (list []models.ConversationInterrupt) {
cnd.Find(db, &list)
return
}
func (r *conversationInterruptRepository) FindOne(db *gorm.DB, cnd *sqls.Cnd) *models.ConversationInterrupt {
ret := &models.ConversationInterrupt{}
if err := cnd.FindOne(db, &ret); err != nil {
return nil
}
return ret
}
func (r *conversationInterruptRepository) Create(db *gorm.DB, t *models.ConversationInterrupt) error {
return db.Create(t).Error
}
func (r *conversationInterruptRepository) Update(db *gorm.DB, t *models.ConversationInterrupt) error {
return db.Save(t).Error
}
func (r *conversationInterruptRepository) Updates(db *gorm.DB, id int64, columns map[string]any) error {
return db.Model(&models.ConversationInterrupt{}).Where("id = ?", id).Updates(columns).Error
}
func (r *conversationInterruptRepository) UpsertByCheckPointID(db *gorm.DB, item *models.ConversationInterrupt) error {
current := r.GetByCheckPointID(db, item.CheckPointID)
if current == nil {
return r.Create(db, item)
}
columns := map[string]any{
"conversation_id": item.ConversationID,
"ai_agent_id": item.AIAgentID,
"agent_run_id": item.AgentRunID,
"agent_step_id": item.AgentStepID,
"source_message_id": item.SourceMessageID,
"last_resume_message_id": item.LastResumeMessageID,
"interrupt_id": item.InterruptID,
"interrupt_type": item.InterruptType,
"status": item.Status,
"prompt_text": item.PromptText,
"request_data": item.RequestData,
"check_point_data": item.CheckPointData,
"resume_count": item.ResumeCount,
"expires_at": item.ExpiresAt,
"updated_at": item.UpdatedAt,
}
return r.Updates(db, current.ID, columns)
}