feat(业务用户组): AUG26-003 业务用户组与店铺负责人分组导入
Some checks failed
构建并部署到测试环境(无 SSH) / build-and-deploy (push) Failing after 1h43m42s
Some checks failed
构建并部署到测试环境(无 SSH) / build-and-deploy (push) Failing after 1h43m42s
- 迁移 000221:新增 tb_business_user_group、tb_business_user_group_member、tb_shop_business_owner_import_task,成员一账号一行由部分唯一索引保证,店铺所属组按当前负责人实时推导,不回填历史分组。 - 用户组 CRUD、成员改组/清空归属、店铺批量交接(原子失败不部分写入)。 - 店铺负责人 CSV 导入任务:逐行独立事务、逐行明细、任务级与行级失败分离。 - 读侧推导与筛选:未分组、业务线、停用组可筛出并带停用标记。 - 补齐操作审计动作与资源、openapi 清单、发布门禁巡检表清单。 - 归档 add-shop-salesperson-groups 变更并同步 openspec/specs/business-user-group,补齐 AUG26-003 验证证据链。
This commit is contained in:
403
internal/task/shop_business_owner_import.go
Normal file
403
internal/task/shop_business_owner_import.go
Normal file
@@ -0,0 +1,403 @@
|
||||
package task
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
"encoding/csv"
|
||||
stderrors "errors"
|
||||
"fmt"
|
||||
"io"
|
||||
"os"
|
||||
"strconv"
|
||||
"strings"
|
||||
|
||||
"github.com/bytedance/sonic"
|
||||
"github.com/hibiken/asynq"
|
||||
"go.uber.org/zap"
|
||||
"gorm.io/gorm"
|
||||
"gorm.io/gorm/clause"
|
||||
|
||||
"github.com/break/junhong_cmp_fiber/internal/infrastructure/audit"
|
||||
"github.com/break/junhong_cmp_fiber/internal/model"
|
||||
"github.com/break/junhong_cmp_fiber/internal/store/postgres"
|
||||
"github.com/break/junhong_cmp_fiber/pkg/auditcontext"
|
||||
"github.com/break/junhong_cmp_fiber/pkg/constants"
|
||||
pkgerrors "github.com/break/junhong_cmp_fiber/pkg/errors"
|
||||
"github.com/break/junhong_cmp_fiber/pkg/storage"
|
||||
"github.com/break/junhong_cmp_fiber/pkg/utils"
|
||||
)
|
||||
|
||||
// ShopBusinessOwnerImportPayload 店铺负责人 CSV 导入任务载荷。
|
||||
type ShopBusinessOwnerImportPayload struct {
|
||||
TaskID uint `json:"task_id"`
|
||||
}
|
||||
|
||||
// ShopBusinessOwnerImportHandler 店铺负责人 CSV 导入任务处理器。
|
||||
// 逐行独立事务:成功行提交、失败行不写店铺并保留原值;任务级失败与行级失败分开记录。
|
||||
type ShopBusinessOwnerImportHandler struct {
|
||||
db *gorm.DB
|
||||
taskStore *postgres.ShopBusinessOwnerImportTaskStore
|
||||
storageService *storage.Service
|
||||
auditWriter *audit.Writer
|
||||
logger *zap.Logger
|
||||
}
|
||||
|
||||
// NewShopBusinessOwnerImportHandler 创建店铺负责人 CSV 导入任务处理器。
|
||||
func NewShopBusinessOwnerImportHandler(
|
||||
db *gorm.DB,
|
||||
taskStore *postgres.ShopBusinessOwnerImportTaskStore,
|
||||
storageService *storage.Service,
|
||||
logger *zap.Logger,
|
||||
auditWriters ...*audit.Writer,
|
||||
) *ShopBusinessOwnerImportHandler {
|
||||
handler := &ShopBusinessOwnerImportHandler{
|
||||
db: db, taskStore: taskStore, storageService: storageService, logger: logger,
|
||||
}
|
||||
if len(auditWriters) > 0 {
|
||||
handler.auditWriter = auditWriters[0]
|
||||
}
|
||||
return handler
|
||||
}
|
||||
|
||||
// Handle 处理店铺负责人 CSV 导入任务。
|
||||
func (h *ShopBusinessOwnerImportHandler) Handle(ctx context.Context, taskMessage *asynq.Task) error {
|
||||
var payload ShopBusinessOwnerImportPayload
|
||||
if err := sonic.Unmarshal(taskMessage.Payload(), &payload); err != nil {
|
||||
h.logger.Error("解析店铺负责人导入任务载荷失败", zap.Error(err))
|
||||
return asynq.SkipRetry
|
||||
}
|
||||
taskRecord, err := h.taskStore.GetByID(ctx, payload.TaskID)
|
||||
if err != nil {
|
||||
h.logger.Error("查询店铺负责人导入任务失败", zap.Uint("task_id", payload.TaskID), zap.Error(err))
|
||||
return asynq.SkipRetry
|
||||
}
|
||||
if h.auditWriter == nil {
|
||||
return pkgerrors.New(pkgerrors.CodeInvalidStatus, "店铺负责人导入统一审计接缝未配置")
|
||||
}
|
||||
ctx = auditcontext.With(ctx, auditcontext.Context{
|
||||
ActorKind: constants.AuditActorSystemTask, ActorID: constants.TaskTypeShopBusinessOwnerImport,
|
||||
ActorName: "店铺负责人导入任务", Source: constants.AuditSourceWorker,
|
||||
CorrelationID: taskRecord.TaskNo,
|
||||
ParentEventID: audit.TaskEventID(constants.AuditResourceShopBusinessOwnerImportTask, taskRecord.ID, "completed"),
|
||||
})
|
||||
claimed, err := h.taskStore.ResetForProcessing(ctx, taskRecord.ID)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if !claimed {
|
||||
h.logger.Info("店铺负责人导入任务已终结,跳过重复消费", zap.Uint("task_id", taskRecord.ID))
|
||||
return nil
|
||||
}
|
||||
|
||||
rows, err := h.downloadAndParse(ctx, taskRecord.StorageKey)
|
||||
if err != nil {
|
||||
h.logger.Warn("下载或解析店铺负责人导入CSV失败", zap.Uint("task_id", taskRecord.ID), zap.Error(err))
|
||||
if finishErr := h.finishTask(ctx, taskRecord, nil, 0, 0, model.ImportTaskStatusFailed, err.Error()); finishErr != nil {
|
||||
return finishErr
|
||||
}
|
||||
return asynq.SkipRetry
|
||||
}
|
||||
|
||||
items, successCount, err := h.processRows(ctx, taskRecord, rows)
|
||||
if err != nil {
|
||||
message := "导入执行中断:" + err.Error()
|
||||
h.logger.Error("店铺负责人导入行执行中断", zap.Uint("task_id", taskRecord.ID), zap.Error(err))
|
||||
if finishErr := h.finishTask(ctx, taskRecord, nil, 0, 0, model.ImportTaskStatusFailed, message); finishErr != nil {
|
||||
return finishErr
|
||||
}
|
||||
return asynq.SkipRetry
|
||||
}
|
||||
failCount := len(items) - successCount
|
||||
if err := h.finishTask(ctx, taskRecord, items, successCount, failCount, model.ImportTaskStatusCompleted, ""); err != nil {
|
||||
return err
|
||||
}
|
||||
h.logger.Info("店铺负责人导入任务完成",
|
||||
zap.Uint("task_id", taskRecord.ID), zap.Int("success", successCount), zap.Int("fail", failCount))
|
||||
return nil
|
||||
}
|
||||
|
||||
// finishTask 在单事务内写任务终态、逐行明细与任务根审计事件。
|
||||
// 任务级失败不产生行明细,与行级失败原因分开记录。
|
||||
func (h *ShopBusinessOwnerImportHandler) finishTask(
|
||||
ctx context.Context,
|
||||
taskRecord *model.ShopBusinessOwnerImportTask,
|
||||
items model.ShopBusinessOwnerImportResults,
|
||||
successCount, failCount, status int,
|
||||
errorMessage string,
|
||||
) error {
|
||||
return h.db.WithContext(ctx).Transaction(func(tx *gorm.DB) error {
|
||||
store := h.taskStore.WithTx(tx)
|
||||
if status == model.ImportTaskStatusFailed {
|
||||
// 任务级失败:未命中非终态说明任务已被其他执行路径终结,按库内事实跳过重复收尾。
|
||||
hit, err := store.MarkFailed(ctx, taskRecord.ID, errorMessage)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if !hit {
|
||||
return nil
|
||||
}
|
||||
} else if err := store.Complete(ctx, taskRecord.ID, len(items), successCount, failCount, items); err != nil {
|
||||
return err
|
||||
}
|
||||
result := batchAuditResult(successCount, failCount)
|
||||
afterData := map[string]any{
|
||||
"status": status, "total_count": len(items), "success_count": successCount, "fail_count": failCount,
|
||||
}
|
||||
if status == model.ImportTaskStatusFailed {
|
||||
result = constants.AuditResultFailed
|
||||
afterData["error_message"] = errorMessage
|
||||
}
|
||||
return h.auditWriter.WriteTask(ctx, tx, audit.TaskInput{
|
||||
EventID: audit.TaskEventID(constants.AuditResourceShopBusinessOwnerImportTask, taskRecord.ID, "completed"),
|
||||
ActionCode: constants.AuditActionShopBusinessOwnerImportTaskCompleted,
|
||||
Summary: "完成店铺负责人导入任务", TaskID: taskRecord.ID, TaskNo: taskRecord.TaskNo,
|
||||
Result: result, CorrelationID: taskRecord.TaskNo,
|
||||
ParentEventID: audit.TaskEventID(constants.AuditResourceShopBusinessOwnerImportTask, taskRecord.ID, "created"),
|
||||
BatchTotal: len(items), SuccessCount: successCount, FailCount: failCount,
|
||||
IdentitySnapshot: map[string]any{
|
||||
"id": taskRecord.ID, "task_no": taskRecord.TaskNo, "file_name": taskRecord.FileName,
|
||||
},
|
||||
BeforeData: map[string]any{"status": model.ImportTaskStatusProcessing},
|
||||
AfterData: afterData,
|
||||
ErrorSummary: errorMessage,
|
||||
})
|
||||
})
|
||||
}
|
||||
|
||||
// downloadAndParse 下载并解析导入 CSV,编码、表头或格式问题一律按任务级失败返回。
|
||||
func (h *ShopBusinessOwnerImportHandler) downloadAndParse(ctx context.Context, key string) ([]shopBusinessOwnerImportRow, error) {
|
||||
if h.storageService == nil {
|
||||
return nil, shopBusinessOwnerImportError("对象存储服务未配置")
|
||||
}
|
||||
if key == "" {
|
||||
return nil, shopBusinessOwnerImportError("导入文件Key不能为空")
|
||||
}
|
||||
localPath, cleanup, err := h.storageService.DownloadToTemp(ctx, key)
|
||||
if err != nil {
|
||||
return nil, shopBusinessOwnerImportError("下载导入CSV失败")
|
||||
}
|
||||
defer cleanup()
|
||||
// 不设行数与体积硬上限:体积沿用上传用途的既有校验,此处按文件实际大小读取。
|
||||
data, err := os.ReadFile(localPath)
|
||||
if err != nil {
|
||||
return nil, shopBusinessOwnerImportError("读取导入CSV失败")
|
||||
}
|
||||
decoded, err := utils.DecodeTextToUTF8(data)
|
||||
if err != nil {
|
||||
return nil, shopBusinessOwnerImportError(constants.ShopBusinessOwnerImportErrorEncoding)
|
||||
}
|
||||
return parseShopBusinessOwnerImportCSV(decoded)
|
||||
}
|
||||
|
||||
// shopBusinessOwnerImportRow 是导入文件的单行业务事实;行号自数据首行起计,表头不计入。
|
||||
type shopBusinessOwnerImportRow struct {
|
||||
Line int
|
||||
ColumnCountMatched bool
|
||||
ShopCode string
|
||||
OperationType string
|
||||
OwnerUsername string
|
||||
Remark string
|
||||
}
|
||||
|
||||
// parseShopBusinessOwnerImportCSV 解析固定列序的导入 CSV。
|
||||
// 表头必须与固定列序完全一致,不一致即任务级失败且不进入逐行阶段;
|
||||
// 数据行列数不符属行级「行格式错误」,因此必须关闭字段数一致性校验,
|
||||
// 否则标准库在首条记录定型字段数后会让后续异常行直接返回 ErrFieldCount,
|
||||
// 把行级问题误升级为任务级失败且不产生行明细。
|
||||
func parseShopBusinessOwnerImportCSV(data []byte) ([]shopBusinessOwnerImportRow, error) {
|
||||
reader := csv.NewReader(bytes.NewReader(data))
|
||||
reader.TrimLeadingSpace = true
|
||||
reader.FieldsPerRecord = -1
|
||||
rows := make([]shopBusinessOwnerImportRow, 0)
|
||||
line := 0
|
||||
for {
|
||||
record, err := reader.Read()
|
||||
if err == io.EOF {
|
||||
break
|
||||
}
|
||||
if err != nil {
|
||||
// 关闭字段数校验后仍报错,说明是引号未闭合等真实 CSV 语法错误,属任务级失败。
|
||||
return nil, shopBusinessOwnerImportError(constants.ShopBusinessOwnerImportErrorFileFormat)
|
||||
}
|
||||
if line == 0 {
|
||||
if !matchShopBusinessOwnerImportHeader(record) {
|
||||
return nil, shopBusinessOwnerImportError(constants.ShopBusinessOwnerImportErrorFileFormat)
|
||||
}
|
||||
line++
|
||||
continue
|
||||
}
|
||||
line++
|
||||
row := shopBusinessOwnerImportRow{Line: line - 1}
|
||||
if len(record) != len(constants.ShopBusinessOwnerImportColumns) {
|
||||
// 列数不符的行不参与业务校验,直接以行格式错误记录并保留原值。
|
||||
rows = append(rows, row)
|
||||
continue
|
||||
}
|
||||
row.ColumnCountMatched = true
|
||||
row.ShopCode = strings.TrimSpace(record[0])
|
||||
row.OperationType = strings.TrimSpace(record[1])
|
||||
row.OwnerUsername = strings.TrimSpace(record[2])
|
||||
row.Remark = strings.TrimSpace(record[3])
|
||||
rows = append(rows, row)
|
||||
}
|
||||
if len(rows) == 0 {
|
||||
return nil, shopBusinessOwnerImportError(constants.ShopBusinessOwnerImportErrorNoDataRow)
|
||||
}
|
||||
return rows, nil
|
||||
}
|
||||
|
||||
// matchShopBusinessOwnerImportHeader 逐列比较表头与固定列序,仅容忍列内两侧空白差异。
|
||||
func matchShopBusinessOwnerImportHeader(record []string) bool {
|
||||
columns := constants.ShopBusinessOwnerImportColumns
|
||||
if len(record) != len(columns) {
|
||||
return false
|
||||
}
|
||||
for index, column := range columns {
|
||||
if strings.TrimSpace(record[index]) != column {
|
||||
return false
|
||||
}
|
||||
}
|
||||
return true
|
||||
}
|
||||
|
||||
// processRows 逐行独立执行并按批更新进度计数;进度写失败不回滚已提交行。
|
||||
// 返回错误表示行执行遇到基础设施故障,由调用方按任务级失败收尾。
|
||||
func (h *ShopBusinessOwnerImportHandler) processRows(ctx context.Context, taskRecord *model.ShopBusinessOwnerImportTask, rows []shopBusinessOwnerImportRow) (model.ShopBusinessOwnerImportResults, int, error) {
|
||||
items := make(model.ShopBusinessOwnerImportResults, 0, len(rows))
|
||||
successCount, failCount := 0, 0
|
||||
for index, row := range rows {
|
||||
item, err := h.processRow(ctx, taskRecord, row)
|
||||
if err != nil {
|
||||
return nil, 0, err
|
||||
}
|
||||
items = append(items, item)
|
||||
if item.Status == constants.ShopBusinessOwnerImportItemStatusSuccess {
|
||||
successCount++
|
||||
} else {
|
||||
failCount++
|
||||
}
|
||||
if (index+1)%constants.ShopBusinessOwnerImportProgressBatchSize == 0 {
|
||||
if err := h.taskStore.UpdateProgress(ctx, taskRecord.ID, len(rows), successCount, failCount); err != nil {
|
||||
h.logger.Warn("更新店铺负责人导入进度失败", zap.Uint("task_id", taskRecord.ID), zap.Error(err))
|
||||
}
|
||||
}
|
||||
}
|
||||
return items, successCount, nil
|
||||
}
|
||||
|
||||
// processRow 校验并执行单行;失败行只记录固定枚举原因,不改动店铺负责人原值。
|
||||
func (h *ShopBusinessOwnerImportHandler) processRow(ctx context.Context, taskRecord *model.ShopBusinessOwnerImportTask, row shopBusinessOwnerImportRow) (model.ShopBusinessOwnerImportResultItem, error) {
|
||||
item := model.ShopBusinessOwnerImportResultItem{Line: row.Line, ShopCode: row.ShopCode, OperationType: row.OperationType}
|
||||
if !row.ColumnCountMatched {
|
||||
return failedShopBusinessOwnerImportItem(item, constants.ShopBusinessOwnerImportRowErrorFormat), nil
|
||||
}
|
||||
var shop model.Shop
|
||||
err := h.db.WithContext(ctx).Select("id", "shop_code", "shop_name", "parent_id", "level", "business_owner_account_id").
|
||||
Where("shop_code = ?", row.ShopCode).First(&shop).Error
|
||||
if err != nil {
|
||||
if err != gorm.ErrRecordNotFound {
|
||||
return item, pkgerrors.Wrap(pkgerrors.CodeDatabaseError, err, "查询导入目标店铺失败")
|
||||
}
|
||||
return failedShopBusinessOwnerImportItem(item, constants.ShopBusinessOwnerImportRowErrorShopMissing), nil
|
||||
}
|
||||
clear := false
|
||||
switch row.OperationType {
|
||||
case constants.ShopBusinessOwnerImportOperationRebind:
|
||||
if row.OwnerUsername == "" {
|
||||
return failedShopBusinessOwnerImportItem(item, constants.ShopBusinessOwnerImportRowErrorOwnerRequired), nil
|
||||
}
|
||||
case constants.ShopBusinessOwnerImportOperationClear:
|
||||
if row.OwnerUsername != "" {
|
||||
return failedShopBusinessOwnerImportItem(item, constants.ShopBusinessOwnerImportRowErrorOwnerForbidden), nil
|
||||
}
|
||||
clear = true
|
||||
default:
|
||||
return failedShopBusinessOwnerImportItem(item, constants.ShopBusinessOwnerImportRowErrorOperation), nil
|
||||
}
|
||||
var ownerID *uint
|
||||
if !clear {
|
||||
var account model.Account
|
||||
if err := h.db.WithContext(ctx).
|
||||
Where("username = ? AND user_type = ? AND status = ?", row.OwnerUsername, constants.UserTypePlatform, constants.StatusEnabled).
|
||||
First(&account).Error; err != nil {
|
||||
if err != gorm.ErrRecordNotFound {
|
||||
return item, pkgerrors.Wrap(pkgerrors.CodeDatabaseError, err, "查询导入目标业务员失败")
|
||||
}
|
||||
return failedShopBusinessOwnerImportItem(item, constants.ShopBusinessOwnerImportRowErrorOwnerInvalid), nil
|
||||
}
|
||||
value := account.ID
|
||||
ownerID = &value
|
||||
}
|
||||
// 每行独立事务:成功行提交,失败行回滚并保留原值。
|
||||
rowErr := h.db.WithContext(ctx).Transaction(func(tx *gorm.DB) error {
|
||||
var locked model.Shop
|
||||
if err := tx.Clauses(clause.Locking{Strength: "UPDATE"}).
|
||||
Select("id", "shop_code", "shop_name", "parent_id", "level", "business_owner_account_id").
|
||||
Where("id = ?", shop.ID).First(&locked).Error; err != nil {
|
||||
if err == gorm.ErrRecordNotFound {
|
||||
return pkgerrors.New(pkgerrors.CodeNotFound, constants.ShopBusinessOwnerImportRowErrorShopMissing)
|
||||
}
|
||||
return pkgerrors.Wrap(pkgerrors.CodeDatabaseError, err, "锁定导入目标店铺失败")
|
||||
}
|
||||
before := locked.BusinessOwnerAccountID
|
||||
if err := tx.Model(&model.Shop{}).Where("id = ?", locked.ID).
|
||||
Updates(map[string]any{"business_owner_account_id": ownerID, "updater": taskRecord.Creator}).Error; err != nil {
|
||||
return pkgerrors.Wrap(pkgerrors.CodeDatabaseError, err, "更新店铺负责人失败")
|
||||
}
|
||||
return h.appendRowAudit(ctx, tx, taskRecord, &locked, before, ownerID, row)
|
||||
})
|
||||
if rowErr != nil {
|
||||
var appErr *pkgerrors.AppError
|
||||
if stderrors.As(rowErr, &appErr) && appErr.Code == pkgerrors.CodeNotFound {
|
||||
return failedShopBusinessOwnerImportItem(item, appErr.Message), nil
|
||||
}
|
||||
return item, rowErr
|
||||
}
|
||||
item.Status = constants.ShopBusinessOwnerImportItemStatusSuccess
|
||||
return item, nil
|
||||
}
|
||||
|
||||
// appendRowAudit 在行事务内写实际变更审计,含负责人前后值与行备注。
|
||||
// 行备注写入事件 Metadata,不进入资源前后值字段。
|
||||
func (h *ShopBusinessOwnerImportHandler) appendRowAudit(ctx context.Context, tx *gorm.DB, taskRecord *model.ShopBusinessOwnerImportTask, shop *model.Shop, before, after *uint, row shopBusinessOwnerImportRow) error {
|
||||
if h.auditWriter == nil {
|
||||
return pkgerrors.New(pkgerrors.CodeInvalidStatus, "店铺负责人导入统一审计接缝未配置")
|
||||
}
|
||||
metadata := map[string]any{
|
||||
"import_task_id": taskRecord.ID, "import_task_no": taskRecord.TaskNo,
|
||||
"line": row.Line, "operation_type": row.OperationType,
|
||||
}
|
||||
if row.Remark != "" {
|
||||
metadata["remark"] = row.Remark
|
||||
}
|
||||
shopID := strconv.FormatUint(uint64(shop.ID), 10)
|
||||
// 稳定事件 ID 由任务与行号决定,任务重复消费时同行为幂等重放。
|
||||
return h.auditWriter.Append(ctx, tx, audit.AppendInput{
|
||||
EventID: audit.TaskEventID(constants.AuditResourceShopBusinessOwnerImportTask, taskRecord.ID, fmt.Sprintf("item:%d", row.Line)),
|
||||
ActionCode: constants.AuditActionShopBusinessOwnerImported, Summary: "导入更新店铺负责人归属",
|
||||
ScopeType: constants.AuditScopeShop, ScopeID: shopID,
|
||||
Result: constants.AuditResultSuccess, Metadata: metadata,
|
||||
Resources: []audit.ResourceInput{{
|
||||
Type: constants.AuditResourceShop, ID: &shopID,
|
||||
Key: shop.ShopCode, DisplayName: shop.ShopName,
|
||||
Relation: constants.AuditResourceRelationPrimary, Role: constants.AuditResourceRoleShopTarget,
|
||||
IdentitySnapshot: map[string]any{
|
||||
"id": shop.ID, "shop_code": shop.ShopCode, "shop_name": shop.ShopName,
|
||||
"parent_id": shop.ParentID, "level": shop.Level,
|
||||
},
|
||||
BeforeData: map[string]any{"business_owner_account_id": before},
|
||||
AfterData: map[string]any{"business_owner_account_id": after},
|
||||
}},
|
||||
})
|
||||
}
|
||||
|
||||
func failedShopBusinessOwnerImportItem(item model.ShopBusinessOwnerImportResultItem, reason string) model.ShopBusinessOwnerImportResultItem {
|
||||
item.Status, item.Reason = constants.ShopBusinessOwnerImportItemStatusFailed, reason
|
||||
return item
|
||||
}
|
||||
|
||||
// shopBusinessOwnerImportError 是任务级失败原因,与行级失败原因分开记录。
|
||||
type shopBusinessOwnerImportError string
|
||||
|
||||
// Error 返回任务级失败原因原文。
|
||||
func (e shopBusinessOwnerImportError) Error() string { return string(e) }
|
||||
Reference in New Issue
Block a user