Files
ai-agent/workflow/dao/session/exec_workflow_result_dao.go
T

105 lines
3.4 KiB
Go

package session
import (
"ai-agent/workflow/consts/public"
sessionDto "ai-agent/workflow/model/dto/session"
"ai-agent/workflow/model/entity"
"context"
"gitea.redpowerfuture.com/red-future/common/beans"
"gitea.redpowerfuture.com/red-future/common/db/gfdb"
"github.com/gogf/gf/v2/util/gconv"
)
var ExecWorkflowResultDao = &execWorkflowResultDao{}
type execWorkflowResultDao struct{}
func (d *execWorkflowResultDao) Insert(ctx context.Context, req *sessionDto.CreateWorkflowResultReq) (id int64, err error) {
var s = new(entity.ExecWorkflowResult)
if err = gconv.Struct(req, &s); err != nil {
return
}
r, err := gfdb.DB(ctx, public.DbNameBlackDeacon).Model(ctx, public.TableNameExecWorkflowResult).Insert(s)
if err != nil {
return
}
return r.LastInsertId()
}
func (d *execWorkflowResultDao) BatchInsert(ctx context.Context, req []*sessionDto.CreateWorkflowResultReq) (rows int64, err error) {
var res []*entity.ExecWorkflowResult
if err = gconv.Structs(req, &res); err != nil {
return
}
r, err := gfdb.DB(ctx, public.DbNameBlackDeacon).Model(ctx, public.TableNameExecWorkflowResult).Data(res).Save()
if err != nil {
return
}
return r.RowsAffected()
}
func (d *execWorkflowResultDao) Delete(ctx context.Context, req *sessionDto.DeleteWorkflowResultReq) (rows int64, err error) {
r, err := gfdb.DB(ctx, public.DbNameBlackDeacon).Model(ctx, public.TableNameExecWorkflowResult).Where(entity.ExecWorkflowResultCol.Id, req.Id).Delete()
if err != nil {
return
}
return r.RowsAffected()
}
func (d *execWorkflowResultDao) List(ctx context.Context, creator string, page *beans.Page) (res []*entity.ExecWorkflowResult, total int, err error) {
m := gfdb.DB(ctx, public.DbNameBlackDeacon).Model(ctx, public.TableNameExecWorkflowResult).
Where(entity.ExecWorkflowResultCol.Creator, creator)
m.OrderDesc(entity.ExecWorkflowResultCol.CreatedAt)
if page != nil {
m.Page(int(page.PageNum), int(page.PageSize))
}
r, total, err := m.AllAndCount(false)
if err != nil {
return
}
err = r.Structs(&res)
return
}
// ListBySession 查询会话下工作流结果(按创建时间倒序)
func (d *execWorkflowResultDao) ListBySession(ctx context.Context, sessionId string) (res []*entity.ExecWorkflowResult, err error) {
r, err := gfdb.DB(ctx, public.DbNameBlackDeacon).Model(ctx, public.TableNameExecWorkflowResult).
Where(entity.ExecWorkflowResultCol.SessionId, sessionId).
OrderDesc(entity.ExecWorkflowResultCol.CreatedAt).
All()
if err != nil {
return
}
err = r.Structs(&res)
return
}
// ListByExecId 查询指定工作流执行记录下的结果文件路径
func (d *execWorkflowResultDao) ListByExecId(ctx context.Context, execId int64) (res []*entity.ExecWorkflowResult, err error) {
r, err := gfdb.DB(ctx, public.DbNameBlackDeacon).Model(ctx, public.TableNameExecWorkflowResult).
Where(entity.ExecWorkflowResultCol.ExecId, execId).
All()
if err != nil {
return
}
err = r.Structs(&res)
return
}
// ListByExecIds 批量查询多个执行记录下的结果文件(按创建时间正序)
func (d *execWorkflowResultDao) ListByExecIds(ctx context.Context, execIds []int64) (res []*entity.ExecWorkflowResult, err error) {
if len(execIds) == 0 {
return
}
r, err := gfdb.DB(ctx, public.DbNameBlackDeacon).Model(ctx, public.TableNameExecWorkflowResult).
WhereIn(entity.ExecWorkflowResultCol.ExecId, execIds).
OrderAsc(entity.ExecWorkflowResultCol.CreatedAt).
All()
if err != nil {
return
}
err = r.Structs(&res)
return
}