package flow import ( flowDao "ai-agent/workflow/dao/flow" flowDto "ai-agent/workflow/model/dto/flow" "ai-agent/workflow/model/entity" "context" "fmt" "sort" "strconv" "gitea.redpowerfuture.com/red-future/common/oss" "gitea.redpowerfuture.com/red-future/common/utils" "github.com/gogf/gf/v2/frame/g" "github.com/gogf/gf/v2/os/gtime" "github.com/gogf/gf/v2/util/gconv" ) var FlowExecutionService = &flowExecutionService{} type flowExecutionService struct{} func (s *flowExecutionService) Get(ctx context.Context, req *flowDto.GetFlowExecutionReq) (res *flowDto.VOFlowExecution, err error) { r, err := flowDao.FlowExecutionDao.Get(ctx, req) if err != nil { return nil, err } res = new(flowDto.VOFlowExecution) res.ImgAddressPrefix, err = oss.GetFileAddressPrefix(ctx) if err != nil { return nil, err } err = gconv.Struct(r, &res) return res, err } func (s *flowExecutionService) List(ctx context.Context, req *flowDto.ListFlowExecutionReq) (res *flowDto.ListFlowExecutionTreeRes, err error) { user, err := utils.GetUserInfo(ctx) if err != nil { return } req.Creator = user.UserName list, _, err := flowDao.FlowExecutionDao.List(ctx, req) if err != nil { return nil, err } executionNumber := make(map[int64]int) var validList []*entity.FlowExecution for _, execution := range list { if g.IsEmpty(execution.OutputParams) { continue } validList = append(validList, execution) } totalValid := len(validList) for idx, execution := range validList { executionNumber[execution.Id] = totalValid - idx } type flowWrap struct { flowNode flowDto.FlowNode createdAt *gtime.Time } dateMap := make(map[string]*[]flowWrap) for _, execution := range validList { createDate := execution.CreatedAt.Format("Y-m-d") flowName := execution.FlowName outputParams := execution.OutputParams num := executionNumber[execution.Id] displayFlowName := fmt.Sprintf("会话-%d(%s)", num, flowName) var tempItems []flowDto.OutputItem for _, paramMap := range outputParams { for tsKey, value := range paramMap { if _, err := strconv.ParseInt(tsKey, 10, 64); err != nil { continue } tempItems = append(tempItems, flowDto.OutputItem{ Content: gconv.String(value), }) } } if len(tempItems) == 0 { continue } suffixCount := make(map[string]int) for idx := range tempItems { item := &tempItems[idx] val := item.Content suffix := "内容" ext := GetFileTypeByPath(val) if ext == "image" { suffix = "图片" } if ext == "video" { suffix = "视频" } if ext == "audio" { suffix = "音频" } if ext == "text" { suffix = "文案" } if ext == "html" { suffix = "HTML" } suffixCount[suffix]++ item.Type = ext item.Label = fmt.Sprintf("%s_%d", suffix, suffixCount[suffix]) } flowNode := flowDto.FlowNode{ FlowName: displayFlowName, Id: execution.Id, SessionId: gconv.String(execution.SessionId), Items: tempItems, } if dateMap[createDate] == nil { dateMap[createDate] = &[]flowWrap{} } *dateMap[createDate] = append(*dateMap[createDate], flowWrap{ flowNode: flowNode, createdAt: execution.CreatedAt, }) } var tree []flowDto.DateNode for date, wraps := range dateMap { sort.Slice(*wraps, func(i, j int) bool { return (*wraps)[i].createdAt.After((*wraps)[j].createdAt) }) var flowNodes []flowDto.FlowNode for _, w := range *wraps { flowNodes = append(flowNodes, w.flowNode) } if len(flowNodes) == 0 { continue } tree = append(tree, flowDto.DateNode{ CreateDate: date, Flows: flowNodes, }) } sort.Slice(tree, func(i, j int) bool { return tree[i].CreateDate > tree[j].CreateDate }) imgPrefix, err := oss.GetFileAddressPrefix(ctx) return &flowDto.ListFlowExecutionTreeRes{ Tree: tree, ImgAddressPrefix: imgPrefix, }, nil } // HttpNodeCallback http节点回调接口 func (s *flowExecutionService) HttpNodeCallback(ctx context.Context) (err error) { r := g.RequestFromCtx(ctx) taskId := r.Get("task_id").String() Notify(taskId, r) return nil }