package phone_asset_association import ( "context" "path/filepath" "strconv" "strings" "time" "github.com/hibiken/asynq" "gorm.io/gorm" "github.com/break/junhong_cmp_fiber/internal/infrastructure/audit" "github.com/break/junhong_cmp_fiber/internal/model" "github.com/break/junhong_cmp_fiber/internal/model/dto" "github.com/break/junhong_cmp_fiber/internal/store" "github.com/break/junhong_cmp_fiber/internal/store/postgres" "github.com/break/junhong_cmp_fiber/pkg/auditfailure" "github.com/break/junhong_cmp_fiber/pkg/constants" "github.com/break/junhong_cmp_fiber/pkg/errors" "github.com/break/junhong_cmp_fiber/pkg/middleware" "github.com/break/junhong_cmp_fiber/pkg/queue" ) // TaskPayload 手机号—资产关联解绑导入 Worker 结构化载荷,与 Worker 侧载荷保持同一 JSON 契约。 type TaskPayload struct { TaskID uint `json:"task_id"` } // New 创建手机号—资产关联后台服务。 // associationStore 用于解除关联,taskStore 与 queueClient 用于受理 CSV 解绑导入任务。 func New( db *gorm.DB, associationStore *postgres.PhoneAssetAssociationStore, taskStore *postgres.PhoneAssetUnbindImportTaskStore, assetIdentifierStore *postgres.AssetIdentifierStore, iotCardStore *postgres.IotCardStore, deviceStore *postgres.DeviceStore, queueClient *queue.Client, auditWriter *audit.Writer, ) *Service { return &Service{ db: db, associationStore: associationStore, taskStore: taskStore, assetIdentifierStore: assetIdentifierStore, iotCardStore: iotCardStore, deviceStore: deviceStore, queueClient: queueClient, auditWriter: auditWriter, } } // CreateImportTask 创建 CSV 解绑导入任务并在同一事务写入创建审计,随后投递到独立导入队列。 // 解绑原因与二次确认在 DTO 层强制;任务级原因随任务行落库,供 Worker 读取后写入每次解除的失效原因。 func (s *Service) CreateImportTask(ctx context.Context, request *dto.CreatePhoneAssetUnbindImportRequest) (*dto.PhoneAssetUnbindImportTaskResponse, error) { if request == nil { return nil, errors.New(errors.CodeInvalidParam) } if !strings.HasPrefix(request.FileKey, constants.PhoneAssetUnbindImportStoragePrefix+"/") { return nil, errors.New(errors.CodeInvalidParam, "导入文件Key不属于指定上传目录") } if !strings.EqualFold(filepath.Ext(request.FileKey), ".csv") { return nil, errors.New(errors.CodeInvalidParam, "解绑导入文件必须为CSV格式") } userID := middleware.GetUserIDFromContext(ctx) if userID == 0 { return nil, errors.New(errors.CodeUnauthorized) } if s.auditWriter == nil { return nil, errors.New(errors.CodeInvalidStatus, "手机号资产解绑导入统一审计接缝未配置") } taskRecord := &model.PhoneAssetUnbindImportTask{ TaskNo: s.taskStore.GenerateTaskNo(), FileName: filepath.Base(request.FileKey), StorageKey: request.FileKey, UnbindReason: request.Reason, Status: model.ImportTaskStatusPending, ResultItems: model.PhoneAssetUnbindImportResults{}, CreatorName: middleware.GetUsernameFromContext(ctx), BaseModel: model.BaseModel{Creator: userID, Updater: userID}, } if err := s.db.WithContext(ctx).Transaction(func(tx *gorm.DB) error { if err := s.taskStore.WithTx(tx).Create(ctx, taskRecord); err != nil { return err } return s.writeTaskAudit(ctx, tx, taskRecord, constants.AuditResultSuccess, nil, nil, "") }); err != nil { return nil, errors.Wrap(errors.CodeDatabaseError, err, "创建手机号资产解绑导入任务失败") } var enqueueErr error if s.queueClient == nil { enqueueErr = errors.New(errors.CodeTaskQueueError, "手机号资产解绑导入任务队列未配置") } else { enqueueErr = s.queueClient.EnqueueTask(ctx, constants.TaskTypePhoneAssetUnbindImport, TaskPayload{TaskID: taskRecord.ID}, asynq.Queue(constants.QueueForTaskType(constants.TaskTypePhoneAssetUnbindImport)), asynq.Timeout(constants.PhoneAssetUnbindImportTaskTimeout)) } if enqueueErr != nil { message := "解绑导入任务入队失败" secondaryErr := s.db.WithContext(ctx).Transaction(func(tx *gorm.DB) error { before := unbindImportTaskState(taskRecord) hit, err := s.taskStore.WithTx(tx).MarkFailed(ctx, taskRecord.ID, message) if err != nil { return err } if !hit { // 任务已到终态:Enqueue 报错但消息实际已投递且 Worker 已跑完。 // 库内才是事实,绝不回写失败态,也不写失败审计。 return nil } // 失败原因必须回填到内存快照,响应与失败审计才与库内一致。 taskRecord.Status, taskRecord.ErrorMessage = model.ImportTaskStatusFailed, message now := time.Now() taskRecord.CompletedAt = &now return s.writeTaskAudit(ctx, tx, taskRecord, constants.AuditResultFailed, before, unbindImportTaskState(taskRecord), strconv.Itoa(errors.CodeTaskQueueError)) }) if secondaryErr != nil { auditfailure.RecordSecondaryWriteFailure(constants.AuditActionPhoneAssetUnbindImportTaskCreated, taskRecord.TaskNo, "", taskRecord.TaskNo, strconv.Itoa(errors.CodeTaskQueueError), secondaryErr) } else if taskRecord.Status != model.ImportTaskStatusFailed { // 未命中非终态时重新读取任务行,让响应反映库内真实终态。 if stored, err := s.taskStore.GetByID(ctx, taskRecord.ID); err == nil { taskRecord = stored } } } return toUnbindImportTaskResponse(taskRecord), nil } // ListImportTasks 分页查询解绑导入任务。 func (s *Service) ListImportTasks(ctx context.Context, request *dto.ListPhoneAssetUnbindImportRequest) (*dto.PhoneAssetUnbindImportTaskPageResult, error) { if request == nil { return nil, errors.New(errors.CodeInvalidParam) } page, pageSize := request.Page, request.PageSize if page <= 0 { page = constants.DefaultPage } if pageSize <= 0 { pageSize = constants.DefaultPageSize } if pageSize > constants.MaxPageSize { pageSize = constants.MaxPageSize } tasks, total, err := s.taskStore.List(ctx, &store.QueryOptions{Page: page, PageSize: pageSize}, request.Status) if err != nil { return nil, errors.Wrap(errors.CodeDatabaseError, err, "查询手机号资产解绑导入任务失败") } items := make([]*dto.PhoneAssetUnbindImportTaskResponse, 0, len(tasks)) for _, taskRecord := range tasks { items = append(items, toUnbindImportTaskResponse(taskRecord)) } return &dto.PhoneAssetUnbindImportTaskPageResult{Items: items, Total: total, Page: page, Size: pageSize}, nil } // GetImportTask 查询解绑导入任务详情与逐行结果。 // 逐行结果含解绑当时的完整手机号快照:关系已失效,只有快照能事后展示被解绑的手机号。 func (s *Service) GetImportTask(ctx context.Context, id uint) (*dto.PhoneAssetUnbindImportTaskDetailResponse, error) { if id == 0 { return nil, errors.New(errors.CodeInvalidParam) } taskRecord, err := s.taskStore.GetByID(ctx, id) if err != nil { if err == gorm.ErrRecordNotFound { return nil, errors.New(errors.CodeNotFound, "手机号资产解绑导入任务不存在") } return nil, errors.Wrap(errors.CodeDatabaseError, err, "查询手机号资产解绑导入任务失败") } items := make([]dto.PhoneAssetUnbindImportItemResponse, 0, len(taskRecord.ResultItems)) for _, item := range taskRecord.ResultItems { items = append(items, dto.PhoneAssetUnbindImportItemResponse{ Line: item.Line, AssetType: item.AssetType, AssetIdentifier: item.AssetIdentifier, AssetID: item.AssetID, UnboundCount: item.UnboundCount, AssociatedPhones: item.AssociatedPhones, Status: item.Status, StatusName: constants.GetPhoneAssetAssociationImportItemStatusName(item.Status), Reason: item.Reason, }) } return &dto.PhoneAssetUnbindImportTaskDetailResponse{ PhoneAssetUnbindImportTaskResponse: *toUnbindImportTaskResponse(taskRecord), Items: items, }, nil } // writeTaskAudit 在调用方事务内写导入任务创建或入队失败审计。 // 任务资源身份快照只保留注册表允许的最小字段。 func (s *Service) writeTaskAudit(ctx context.Context, tx *gorm.DB, task *model.PhoneAssetUnbindImportTask, result string, before, after map[string]any, errorCode string) error { return s.auditWriter.WriteTask(ctx, tx, audit.TaskInput{ EventID: audit.TaskEventID(constants.AuditResourcePhoneAssetUnbindImportTask, task.ID, unbindImportTaskAuditPhase(result)), ActionCode: constants.AuditActionPhoneAssetUnbindImportTaskCreated, Summary: "创建手机号资产解绑导入任务", TaskID: task.ID, TaskNo: task.TaskNo, Actor: audit.ActorInput{ Kind: constants.AuditActorAccount, ID: strconv.FormatUint(uint64(middleware.GetUserIDFromContext(ctx)), 10), Name: middleware.GetUsernameFromContext(ctx), }, Source: constants.AuditSourceAdminAPI, ScopeType: constants.AuditScopePlatform, Result: result, ErrorCode: errorCode, ErrorSummary: task.ErrorMessage, // 与 Worker 侧完成事件使用同一关联键,使同一任务的全部事件可按 correlation 串成一条时间线。 CorrelationID: task.TaskNo, IdentitySnapshot: map[string]any{ "id": task.ID, "task_no": task.TaskNo, "file_name": task.FileName, }, BeforeData: before, AfterData: after, }) } // unbindImportTaskAuditPhase 返回任务创建阶段的稳定事件阶段名,失败入队使用独立阶段避免覆盖首次事件。 func unbindImportTaskAuditPhase(result string) string { if result == constants.AuditResultSuccess { return "created" } return "enqueue_failed" } // unbindImportTaskState 生成解绑导入任务状态快照。 func unbindImportTaskState(task *model.PhoneAssetUnbindImportTask) map[string]any { if task == nil { return nil } return map[string]any{ "status": task.Status, "total_count": task.TotalCount, "success_count": task.SuccessCount, "fail_count": task.FailCount, } } func toUnbindImportTaskResponse(taskRecord *model.PhoneAssetUnbindImportTask) *dto.PhoneAssetUnbindImportTaskResponse { response := &dto.PhoneAssetUnbindImportTaskResponse{ ID: taskRecord.ID, TaskNo: taskRecord.TaskNo, FileName: taskRecord.FileName, UnbindReason: taskRecord.UnbindReason, Status: taskRecord.Status, StatusName: model.ImportTaskStatusName(taskRecord.Status), TotalCount: taskRecord.TotalCount, SuccessCount: taskRecord.SuccessCount, FailCount: taskRecord.FailCount, ErrorMessage: taskRecord.ErrorMessage, CreatorName: taskRecord.CreatorName, CreatedAt: taskRecord.CreatedAt.Format(time.RFC3339), } if taskRecord.StartedAt != nil { response.StartedAt = taskRecord.StartedAt.Format(time.RFC3339) } if taskRecord.CompletedAt != nil { response.CompletedAt = taskRecord.CompletedAt.Format(time.RFC3339) } return response }