Nessuna descrizione
You can not select more than 25 topics Topics must start with a letter or number, can include dashes ('-') and can be up to 35 characters long.

batch_save_table_alias_flow.go 6.4KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177
  1. package aliasmanagement
  2. import (
  3. "context"
  4. "fmt"
  5. "git.x2erp.com/qdy/go-base/ctx"
  6. "git.x2erp.com/qdy/go-base/logger"
  7. "git.x2erp.com/qdy/go-base/model/response"
  8. "git.x2erp.com/qdy/go-base/util"
  9. "git.x2erp.com/qdy/go-db/factory/database"
  10. "git.x2erp.com/qdy/go-svc-configure/internal/tables"
  11. "github.com/google/uuid"
  12. "github.com/jmoiron/sqlx"
  13. )
  14. // BatchSaveTableAliasFlow 批量保存表别名字典流水记录
  15. func BatchSaveTableAliasFlow(req *BatchTableAliasRequest, tenantID string, ctx context.Context, dbFactory *database.DBFactory, reqCtx *ctx.RequestContext) *response.QueryResult[[]tables.DicTableAliasFlowDB] {
  16. logger.Debug(fmt.Sprintf("BatchSaveTableAliasFlow-开始批量保存表别名字典流水,租户ID: %s, 数量: %d", tenantID, len(req.Items)))
  17. // 参数验证
  18. if tenantID == "" {
  19. logger.ErrorC(reqCtx, "租户ID不能为空")
  20. return util.CreateErrorResult[[]tables.DicTableAliasFlowDB]("租户ID不能为空", reqCtx)
  21. }
  22. if len(req.Items) == 0 {
  23. logger.ErrorC(reqCtx, "批量保存的表别名字典列表不能为空")
  24. return util.CreateErrorResult[[]tables.DicTableAliasFlowDB]("批量保存的表别名字典列表不能为空", reqCtx)
  25. }
  26. for i, item := range req.Items {
  27. if err := validateTableAliasRequest(&item); err != nil {
  28. logger.ErrorC(reqCtx, fmt.Sprintf("第%d个表别名参数验证失败: %v", i+1, err))
  29. return util.CreateErrorResult[[]tables.DicTableAliasFlowDB](fmt.Sprintf("第%d个表别名参数验证失败: %v", i+1, err), reqCtx)
  30. }
  31. }
  32. // 获取数据库连接并开始事务
  33. db := dbFactory.GetDB()
  34. tx, err := db.BeginTxx(ctx, nil)
  35. if err != nil {
  36. logger.ErrorC(reqCtx, fmt.Sprintf("开始事务失败: %v", err))
  37. return util.CreateErrorResult[[]tables.DicTableAliasFlowDB](fmt.Sprintf("开始事务失败: %v", err), reqCtx)
  38. }
  39. defer func() {
  40. if p := recover(); p != nil {
  41. tx.Rollback()
  42. panic(p)
  43. }
  44. }()
  45. var savedItems []tables.DicTableAliasFlowDB
  46. var errors []string
  47. // 批量处理每个表别名流水记录
  48. for i, item := range req.Items {
  49. logger.Debug(fmt.Sprintf("处理第%d个表别名流水记录: table_id=%s, table_alias=%s", i+1, item.TableID, item.TableAlias))
  50. // 检查是否已存在相同的流水记录(基于table_alias和tenant_id)
  51. exists, err := checkTableAliasFlowExists(ctx, tx, item.TableAlias, tenantID)
  52. if err != nil {
  53. errors = append(errors, fmt.Sprintf("第%d个表别名流水记录检查存在性失败: %v", i+1, err))
  54. continue
  55. }
  56. var tableAliasFlow tables.DicTableAliasFlowDB
  57. if exists {
  58. // 更新流水记录(保持审批状态不变)
  59. tableAliasFlow, err = updateTableAliasFlow(ctx, tx, &item, tenantID)
  60. if err != nil {
  61. errors = append(errors, fmt.Sprintf("第%d个表别名流水记录更新失败: %v", i+1, err))
  62. continue
  63. }
  64. logger.Debug(fmt.Sprintf("更新表别名流水记录成功: %s", item.TableAlias))
  65. } else {
  66. // 插入新的流水记录
  67. tableAliasFlow, err = insertTableAliasFlow(ctx, tx, &item, tenantID)
  68. if err != nil {
  69. errors = append(errors, fmt.Sprintf("第%d个表别名流水记录插入失败: %v", i+1, err))
  70. continue
  71. }
  72. logger.Debug(fmt.Sprintf("插入表别名流水记录成功: %s", item.TableAlias))
  73. }
  74. savedItems = append(savedItems, tableAliasFlow)
  75. }
  76. // 如果有错误,回滚事务
  77. if len(errors) > 0 {
  78. tx.Rollback()
  79. errorMsg := ""
  80. for i, errMsg := range errors {
  81. if i > 0 {
  82. errorMsg += "; "
  83. }
  84. errorMsg += errMsg
  85. }
  86. logger.ErrorC(reqCtx, fmt.Sprintf("批量保存表别名字典流水失败: %s", errorMsg))
  87. return util.CreateErrorResult[[]tables.DicTableAliasFlowDB](fmt.Sprintf("批量保存表别名字典流水失败: %s", errorMsg), reqCtx)
  88. }
  89. // 提交事务
  90. if err := tx.Commit(); err != nil {
  91. logger.ErrorC(reqCtx, fmt.Sprintf("提交事务失败: %v", err))
  92. return util.CreateErrorResult[[]tables.DicTableAliasFlowDB](fmt.Sprintf("提交事务失败: %v", err), reqCtx)
  93. }
  94. logger.Debug(fmt.Sprintf("成功批量保存 %d 个表别名字典流水记录", len(savedItems)))
  95. return util.CreateSuccessResultData[[]tables.DicTableAliasFlowDB](savedItems, reqCtx)
  96. }
  97. // checkTableAliasFlowExists 检查表别名流水记录是否存在(基于table_alias和tenant_id)
  98. func checkTableAliasFlowExists(ctx context.Context, tx *sqlx.Tx, tableAlias, tenantID string) (bool, error) {
  99. var count int
  100. query := "SELECT COUNT(*) FROM dic_table_alias_flow WHERE table_alias = ? AND tenant_id = ? AND deleted_at IS NULL"
  101. err := tx.GetContext(ctx, &count, query, tableAlias, tenantID)
  102. return count > 0, err
  103. }
  104. // updateTableAliasFlow 更新表别名流水记录
  105. func updateTableAliasFlow(ctx context.Context, tx *sqlx.Tx, req *TableAliasRequest, tenantID string) (tables.DicTableAliasFlowDB, error) {
  106. query := `
  107. UPDATE dic_table_alias_flow
  108. SET table_id = ?, updated_at = CURRENT_TIMESTAMP
  109. WHERE table_alias = ? AND tenant_id = ? AND deleted_at IS NULL
  110. `
  111. _, err := tx.ExecContext(ctx, query,
  112. req.TableID,
  113. req.TableAlias,
  114. tenantID,
  115. )
  116. if err != nil {
  117. return tables.DicTableAliasFlowDB{}, err
  118. }
  119. // 查询更新后的记录
  120. var tableAliasFlow tables.DicTableAliasFlowDB
  121. selectQuery := `
  122. SELECT id, table_id, table_alias, tenant_id, approval_status, approver, approved_at, created_at, updated_at, deleted_at
  123. FROM dic_table_alias_flow
  124. WHERE table_alias = ? AND tenant_id = ? AND deleted_at IS NULL
  125. `
  126. err = tx.GetContext(ctx, &tableAliasFlow, selectQuery, req.TableAlias, tenantID)
  127. return tableAliasFlow, err
  128. }
  129. // insertTableAliasFlow 插入表别名流水记录
  130. func insertTableAliasFlow(ctx context.Context, tx *sqlx.Tx, req *TableAliasRequest, tenantID string) (tables.DicTableAliasFlowDB, error) {
  131. // 生成ID
  132. id := uuid.New().String()
  133. query := `
  134. INSERT INTO dic_table_alias_flow (id, table_id, table_alias, tenant_id, approval_status, created_at, updated_at)
  135. VALUES (?, ?, ?, ?, 0, CURRENT_TIMESTAMP, CURRENT_TIMESTAMP)
  136. `
  137. _, err := tx.ExecContext(ctx, query,
  138. id,
  139. req.TableID,
  140. req.TableAlias,
  141. tenantID,
  142. )
  143. if err != nil {
  144. return tables.DicTableAliasFlowDB{}, err
  145. }
  146. // 查询刚插入的记录
  147. var tableAliasFlow tables.DicTableAliasFlowDB
  148. selectQuery := `
  149. SELECT id, table_id, table_alias, tenant_id, approval_status, approver, approved_at, created_at, updated_at, deleted_at
  150. FROM dic_table_alias_flow
  151. WHERE table_alias = ? AND tenant_id = ? AND deleted_at IS NULL
  152. `
  153. err = tx.GetContext(ctx, &tableAliasFlow, selectQuery, req.TableAlias, tenantID)
  154. return tableAliasFlow, err
  155. }