diff --git a/.scratch/tech-global-audit/审计覆盖基线.md b/.scratch/tech-global-audit/审计覆盖基线.md index 51b1bc9..1f216c8 100644 --- a/.scratch/tech-global-audit/审计覆盖基线.md +++ b/.scratch/tech-global-audit/审计覆盖基线.md @@ -73,7 +73,7 @@ | 钱包流水与代理充值业务导出 | N/A:只读投影,不改变钱包、充值或审批状态;任务创建操作者与店铺权限快照沿用现有导出任务记录 | N/A:只读取主钱包流水、充值记录和本地通用审批实例;金额及余额使用既有权威事实,不写入领域账本 | N/A:不实时调用支付渠道或企业微信;明文业务凭证 Key 只进入导出结果,不复制到 Integration Log,且仍禁止记录 Secret、access_token、media_id 和附件正文 | 仅沿用现有 `export:dispatch` → `export:shard` → `export:finalize` Asynq 链路,不新增业务 Outbox | | 退款与换货业务导出 | N/A:只读投影,不改变退款、换货、资产或审批状态;任务创建操作者与店铺权限快照沿用现有导出任务记录 | N/A:只读取退款、订单、套餐使用、换货资产快照和本地审批事实;金额及处理标记沿用既有权威事实,不写入领域账本 | N/A:不实时调用企业微信、支付或 Gateway;明文业务凭证和收货资料只进入有权导出结果,不复制到 Integration Log,且仍禁止记录 Secret、access_token、media_id 和附件正文 | 仅沿用现有 `export:dispatch` → `export:shard` → `export:finalize` Asynq 链路,不新增业务 Outbox | | IoT 卡固定档位限速 | 复用资产操作审计 `card_speed_tier`,记录后台操作者、卡 ICCID、固定档位、Integration Log ID 和 success/failed 结果;设备无入口且不通过绑定卡间接限速 | N/A:不在本地保存或修改卡当前限速状态,Gateway 是外部执行方 | 每次实际 Gateway 调用前写 pending,按 success/failed/unknown 终结;超时 unknown 保存按 ICCID 人工核对策略,摘要不含 Secret、access_token 或完整响应正文 | N/A:单次外部命令无后续可靠副作用,结果未知禁止盲目重发,不创建自动补偿 Outbox | -| 设备 CSV 批量分配代理或套餐系列 | 任务记录冻结操作者和目标;逐批复用 `device_allocate` 或 `device_series_binding` 资产操作审计,记录设备前后值、目标、成功/失败数和失败原因 | N/A:设备归属、绑定卡归属、分配记录及 `series_id` 是权威业务事实,不另建领域账本 | N/A:CSV 解析和两类分配均为本地数据库操作,不调用 Gateway、支付或企微;对象存储沿用现有存储日志 | 复用 `device:import` Asynq 任务;状态条件阻止完成任务重复执行,处理中断恢复时已达到目标关系的设备按成功处理,不新增业务 Outbox | +| 设备 CSV 批量分配、设置套餐系列或回收 | 任务记录冻结操作者和可选目标;逐批复用 `device_allocate`、`device_series_binding` 或 `device_recall` 资产操作审计,记录设备前后值、目标、成功/失败数和失败原因 | N/A:设备归属、绑定卡归属、分配记录及 `series_id` 是权威业务事实,不另建领域账本 | N/A:CSV 解析、分配和回收均为本地数据库操作,不调用 Gateway、支付或企微;对象存储沿用现有存储日志 | 复用 `device:import` Asynq 任务;状态条件阻止完成任务重复执行,处理中断恢复时已达到目标关系的设备按成功处理,不新增业务 Outbox | ### `deliver-july-iteration-confirmed-scope` 任务覆盖映射 diff --git a/README.md b/README.md index 3d4cbb7..d7d908d 100644 --- a/README.md +++ b/README.md @@ -954,7 +954,7 @@ rdb.Set(ctx, key, status, time.Hour) - **[资产套餐批量订购](docs/asset-package-batch-order/功能总结.md)**:单列 CSV、整批统一套餐与支付方式、五态异步任务及逐行订单结果说明 - **[业务数据导出](docs/business-data-export/功能总结.md)**:IoT 卡套餐流量列、套餐列表及后续业务 datasource 的字段、筛选和权限口径 - **[IoT 卡固定档位限速](docs/iot-card-speed-tier/功能总结.md)**:仅按卡 ICCID 调用 Gateway 的 `-1..8` 固定档位、权限、Integration Log 与结果未知处理;设备不提供限速入口 -- **[设备 CSV 批量分配](docs/device-batch-allocation/功能总结.md)**:复用设备导入任务外壳,以单列 VirtualNo/IMEI/SN 为整批分配目标代理或套餐系列,并沿用现有权限、事务、审计和幂等规则 +- **[设备 CSV 批量分配与回收](docs/device-batch-allocation/功能总结.md)**:复用设备导入任务外壳,以单列 VirtualNo/IMEI/SN 整批分配目标代理、设置套餐系列或回收设备,并沿用现有权限、事务、审计和幂等规则 - **[快速开始指南](specs/001-fiber-middleware-integration/quickstart.md)**:详细设置和测试说明 - **[限流指南](docs/rate-limiting.md)**:全面的限流配置和使用 - **[错误处理使用指南](docs/003-error-handling/使用指南.md)**:错误码参考、Handler 使用、客户端处理、最佳实践 diff --git a/docs/admin-openapi.yaml b/docs/admin-openapi.yaml index bcf69ee..c5586b2 100644 --- a/docs/admin-openapi.yaml +++ b/docs/admin-openapi.yaml @@ -3784,19 +3784,19 @@ components: minLength: 1 type: string operation_type: - description: 操作类型 (assign_shop:分配目标代理, assign_series:设置套餐系列) + description: 操作类型 (assign_shop:分配目标代理, assign_series:设置套餐系列, recall:回收设备) enum: - assign_shop - assign_series + - recall type: string target_id: - description: 目标代理店铺ID或套餐系列ID,由 operation_type 决定 + description: 目标代理店铺ID或套餐系列ID;recall 时不传 minimum: 1 type: integer required: - file_key - operation_type - - target_id type: object DtoCreateDeviceBatchAllocationResponse: properties: @@ -5051,7 +5051,7 @@ components: description: 任务业务类型名称(中文) type: string operation_type: - description: 任务业务类型 (import:导入设备, assign_shop:分配目标代理, assign_series:设置套餐系列) + description: 任务业务类型 (import:导入设备, assign_shop:分配目标代理, assign_series:设置套餐系列, recall:回收设备) type: string realname_policy: description: 实名认证策略 (none:无需实名, before_order:先实名后充值/购买, after_order:先充值/购买后实名) @@ -5083,7 +5083,7 @@ components: description: 成功数 type: integer target_id: - description: 批量分配目标店铺或套餐系列ID + description: 批量分配目标店铺或套餐系列ID,回收任务为空 minimum: 0 nullable: true type: integer @@ -5137,7 +5137,7 @@ components: description: 任务业务类型名称(中文) type: string operation_type: - description: 任务业务类型 (import:导入设备, assign_shop:分配目标代理, assign_series:设置套餐系列) + description: 任务业务类型 (import:导入设备, assign_shop:分配目标代理, assign_series:设置套餐系列, recall:回收设备) type: string realname_policy: description: 实名认证策略 (none:无需实名, before_order:先实名后充值/购买, after_order:先充值/购买后实名) @@ -5163,7 +5163,7 @@ components: description: 成功数 type: integer target_id: - description: 批量分配目标店铺或套餐系列ID + description: 批量分配目标店铺或套餐系列ID,回收任务为空 minimum: 0 nullable: true type: integer @@ -6296,7 +6296,7 @@ components: minLength: 1 type: string purpose: - description: 文件用途 (iot_import:ICCID导入, export:数据导出, attachment:附件, batch_purchase:资产套餐批量订购CSV, device_batch_allocation:设备批量分配CSV) + description: 文件用途 (iot_import:ICCID导入, export:数据导出, attachment:附件, batch_purchase:资产套餐批量订购CSV, device_batch_allocation:设备批量分配或回收CSV) enum: - iot_import - export @@ -17181,7 +17181,7 @@ paths: - 设备管理 /api/admin/devices/import/allocations: post: - description: 复用设备导入任务。CSV只能包含一列设备标识,每行支持VirtualNo、IMEI或SN;整批选择分配目标代理或设置目标套餐系列。 + description: 复用设备导入任务。CSV只能包含一列设备标识,每行支持VirtualNo、IMEI或SN;整批可分配目标代理、设置目标套餐系列或回收设备。 requestBody: content: application/json: @@ -17240,7 +17240,7 @@ paths: description: 服务器内部错误 security: - BearerAuth: [] - summary: 创建CSV设备批量分配任务 + summary: 创建CSV设备批量分配或回收任务 tags: - 设备管理 /api/admin/devices/import/tasks: @@ -17271,15 +17271,16 @@ paths: minimum: 1 nullable: true type: integer - - description: 任务业务类型 (import:导入设备, assign_shop:分配目标代理, assign_series:设置套餐系列) + - description: 任务业务类型 (import:导入设备, assign_shop:分配目标代理, assign_series:设置套餐系列, recall:回收设备) in: query name: operation_type schema: - description: 任务业务类型 (import:导入设备, assign_shop:分配目标代理, assign_series:设置套餐系列) + description: 任务业务类型 (import:导入设备, assign_shop:分配目标代理, assign_series:设置套餐系列, recall:回收设备) enum: - import - assign_shop - assign_series + - recall type: string - description: 批次号(模糊查询) in: query @@ -27501,7 +27502,7 @@ paths: |---|------|-------------| | iot_import | ICCID/设备导入 (Excel) | imports/YYYY/MM/DD/uuid.xlsx | | batch_purchase | 资产套餐批量订购 (CSV) | batch-purchases/YYYY/MM/DD/uuid.csv | - | device_batch_allocation | 设备批量分配代理或套餐系列 (CSV) | device-batch-allocations/YYYY/MM/DD/uuid.csv | + | device_batch_allocation | 设备批量分配、设置套餐系列或回收 (CSV) | device-batch-allocations/YYYY/MM/DD/uuid.csv | | export | 数据导出 | exports/YYYY/MM/DD/uuid.xlsx | | attachment | 附件上传 | attachments/YYYY/MM/DD/uuid.ext | diff --git a/docs/device-batch-allocation/功能总结.md b/docs/device-batch-allocation/功能总结.md index 75be401..91d5cd0 100644 --- a/docs/device-batch-allocation/功能总结.md +++ b/docs/device-batch-allocation/功能总结.md @@ -1,4 +1,4 @@ -# 设备 CSV 批量分配功能总结 +# 设备 CSV 批量分配与回收功能总结 ## 功能边界 @@ -22,7 +22,8 @@ POST /api/admin/devices/import/allocations - `assign_shop`:`target_id` 为目标代理店铺 ID。 - `assign_series`:`target_id` 为目标套餐系列 ID。 -- 每个任务只能选择一种操作和一个目标,不能逐行指定不同目标。 +- `recall`:不传 `target_id`;平台回收到平台库存,代理回收到自己店铺。 +- 每个任务只能选择一种操作,不能逐行指定不同操作或目标。 ## CSV 格式 @@ -34,8 +35,9 @@ Worker 批量解析设备后调用现有业务方法: - 分配代理调用 `DeviceService.AllocateDevices`,沿用直属下级、平台库存、绑定卡归属同步、分配记录和资产审计规则。 - 设置套餐系列调用 `DeviceService.BatchSetSeriesBinding`,沿用套餐系列有效性、代理授权、设备归属和资产审计规则。 +- 回收设备调用 `DeviceService.RecallDevices`,代理只能回收直属下级的设备,平台可回收店铺设备,并同步绑定卡归属、分配记录和资产审计。 -任务创建时冻结操作者 ID、类型和店铺 ID,Worker 重建权限上下文,不能因异步执行丢失数据范围。重复消费时,已经属于目标店铺或已经绑定目标系列的设备直接按成功处理;任务完成或失败后不重复执行。CSV 内重复设备记录为跳过,不产生重复副作用。 +任务创建时冻结操作者 ID、类型和店铺 ID,Worker 重建权限上下文,不能因异步执行丢失数据范围。重复消费时,已经属于目标店铺、已经绑定目标系列或已经回收到操作者归属的设备直接按成功处理;任务完成或失败后不重复执行。CSV 内重复设备记录为跳过,不产生重复副作用。 ## 任务查询与回滚 diff --git a/internal/handler/admin/device_import.go b/internal/handler/admin/device_import.go index e1e62f6..a748716 100644 --- a/internal/handler/admin/device_import.go +++ b/internal/handler/admin/device_import.go @@ -46,7 +46,7 @@ func (h *DeviceImportHandler) Import(c *fiber.Ctx) error { return response.Success(c, result) } -// CreateAllocation 创建单列 CSV 设备批量分配任务。 +// CreateAllocation 创建单列 CSV 设备批量分配或回收任务。 // POST /api/admin/devices/import/allocations func (h *DeviceImportHandler) CreateAllocation(c *fiber.Ctx) error { var req dto.CreateDeviceBatchAllocationRequest diff --git a/internal/model/device_import_task.go b/internal/model/device_import_task.go index f7de5c5..c058b89 100644 --- a/internal/model/device_import_task.go +++ b/internal/model/device_import_task.go @@ -13,8 +13,8 @@ type DeviceImportTask struct { gorm.Model BaseModel `gorm:"embedded"` TaskNo string `gorm:"column:task_no;type:varchar(50);uniqueIndex:idx_device_import_task_no,where:deleted_at IS NULL;not null;comment:任务编号(唯一)" json:"task_no"` - OperationType string `gorm:"column:operation_type;type:varchar(32);not null;default:'import';comment:任务业务类型(import/assign_shop/assign_series)" json:"operation_type"` - TargetID *uint `gorm:"column:target_id;comment:批量分配目标店铺或套餐系列ID" json:"target_id"` + OperationType string `gorm:"column:operation_type;type:varchar(32);not null;default:'import';comment:任务业务类型(import/assign_shop/assign_series/recall)" json:"operation_type"` + TargetID *uint `gorm:"column:target_id;comment:批量分配目标店铺或套餐系列ID,回收任务为空" json:"target_id"` OperatorType int `gorm:"column:operator_type;type:int;not null;default:0;comment:任务创建时操作者类型快照" json:"operator_type"` OperatorShopID *uint `gorm:"column:operator_shop_id;comment:任务创建时操作者店铺ID快照" json:"operator_shop_id"` BatchNo string `gorm:"column:batch_no;type:varchar(100);comment:批次号" json:"batch_no"` diff --git a/internal/model/dto/device_import_dto.go b/internal/model/dto/device_import_dto.go index 4c9d99e..5644672 100644 --- a/internal/model/dto/device_import_dto.go +++ b/internal/model/dto/device_import_dto.go @@ -14,14 +14,14 @@ type ImportDeviceResponse struct { Message string `json:"message" description:"提示信息"` } -// CreateDeviceBatchAllocationRequest 创建设备 CSV 批量分配任务请求。 +// CreateDeviceBatchAllocationRequest 创建设备 CSV 批量操作任务请求。 type CreateDeviceBatchAllocationRequest struct { FileKey string `json:"file_key" validate:"required,min=1,max=500" required:"true" minLength:"1" maxLength:"500" description:"单列设备标识 CSV 对象存储 Key(每行支持 VirtualNo、IMEI 或 SN)"` - OperationType string `json:"operation_type" validate:"required,oneof=assign_shop assign_series" required:"true" enum:"assign_shop,assign_series" description:"操作类型 (assign_shop:分配目标代理, assign_series:设置套餐系列)"` - TargetID uint `json:"target_id" validate:"required,min=1" required:"true" minimum:"1" description:"目标代理店铺ID或套餐系列ID,由 operation_type 决定"` + OperationType string `json:"operation_type" validate:"required,oneof=assign_shop assign_series recall" required:"true" enum:"assign_shop,assign_series,recall" description:"操作类型 (assign_shop:分配目标代理, assign_series:设置套餐系列, recall:回收设备)"` + TargetID uint `json:"target_id,omitempty" validate:"omitempty,min=1" minimum:"1" description:"目标代理店铺ID或套餐系列ID;recall 时不传"` } -// CreateDeviceBatchAllocationResponse 创建设备 CSV 批量分配任务响应。 +// CreateDeviceBatchAllocationResponse 创建设备 CSV 批量操作任务响应。 type CreateDeviceBatchAllocationResponse struct { TaskID uint `json:"task_id" description:"任务ID"` TaskNo string `json:"task_no" description:"任务编号"` @@ -32,7 +32,7 @@ type ListDeviceImportTaskRequest struct { Page int `json:"page" query:"page" validate:"omitempty,min=1" minimum:"1" description:"页码"` PageSize int `json:"page_size" query:"page_size" validate:"omitempty,min=1,max=100" minimum:"1" maximum:"100" description:"每页数量"` Status *int `json:"status" query:"status" validate:"omitempty,min=1,max=4" minimum:"1" maximum:"4" description:"任务状态 (1:待处理, 2:处理中, 3:已完成, 4:失败)"` - OperationType string `json:"operation_type" query:"operation_type" validate:"omitempty,oneof=import assign_shop assign_series" enum:"import,assign_shop,assign_series" description:"任务业务类型 (import:导入设备, assign_shop:分配目标代理, assign_series:设置套餐系列)"` + OperationType string `json:"operation_type" query:"operation_type" validate:"omitempty,oneof=import assign_shop assign_series recall" enum:"import,assign_shop,assign_series,recall" description:"任务业务类型 (import:导入设备, assign_shop:分配目标代理, assign_series:设置套餐系列, recall:回收设备)"` BatchNo string `json:"batch_no" query:"batch_no" validate:"omitempty,max=100" maxLength:"100" description:"批次号(模糊查询)"` StartTime *time.Time `json:"start_time" query:"start_time" description:"创建时间起始"` EndTime *time.Time `json:"end_time" query:"end_time" description:"创建时间结束"` @@ -41,9 +41,9 @@ type ListDeviceImportTaskRequest struct { type DeviceImportTaskResponse struct { ID uint `json:"id" description:"任务ID"` TaskNo string `json:"task_no" description:"任务编号"` - OperationType string `json:"operation_type" description:"任务业务类型 (import:导入设备, assign_shop:分配目标代理, assign_series:设置套餐系列)"` + OperationType string `json:"operation_type" description:"任务业务类型 (import:导入设备, assign_shop:分配目标代理, assign_series:设置套餐系列, recall:回收设备)"` OperationName string `json:"operation_name" description:"任务业务类型名称(中文)"` - TargetID *uint `json:"target_id,omitempty" description:"批量分配目标店铺或套餐系列ID"` + TargetID *uint `json:"target_id,omitempty" description:"批量分配目标店铺或套餐系列ID,回收任务为空"` Status int `json:"status" description:"任务状态 (1:待处理, 2:处理中, 3:已完成, 4:失败)"` StatusName string `json:"status_name" description:"任务状态名称(中文)"` StatusText string `json:"status_text" description:"任务状态文本"` diff --git a/internal/model/dto/storage_dto.go b/internal/model/dto/storage_dto.go index 46304e3..f783548 100644 --- a/internal/model/dto/storage_dto.go +++ b/internal/model/dto/storage_dto.go @@ -3,7 +3,7 @@ package dto type GetUploadURLRequest struct { FileName string `json:"file_name" validate:"required,min=1,max=255" required:"true" minLength:"1" maxLength:"255" description:"文件名(如:cards.csv)"` ContentType string `json:"content_type" validate:"omitempty,max=100" maxLength:"100" description:"文件 MIME 类型(如:text/csv),留空则自动推断"` - Purpose string `json:"purpose" validate:"required,oneof=iot_import export attachment batch_purchase device_batch_allocation" required:"true" enum:"iot_import,export,attachment,batch_purchase,device_batch_allocation" description:"文件用途 (iot_import:ICCID导入, export:数据导出, attachment:附件, batch_purchase:资产套餐批量订购CSV, device_batch_allocation:设备批量分配CSV)"` + Purpose string `json:"purpose" validate:"required,oneof=iot_import export attachment batch_purchase device_batch_allocation" required:"true" enum:"iot_import,export,attachment,batch_purchase,device_batch_allocation" description:"文件用途 (iot_import:ICCID导入, export:数据导出, attachment:附件, batch_purchase:资产套餐批量订购CSV, device_batch_allocation:设备批量分配或回收CSV)"` } type GetUploadURLResponse struct { diff --git a/internal/routes/device.go b/internal/routes/device.go index 8582775..2fa065e 100644 --- a/internal/routes/device.go +++ b/internal/routes/device.go @@ -125,8 +125,8 @@ func registerDeviceRoutes(router fiber.Router, handler *admin.DeviceHandler, imp }) Register(devices, doc, groupPath, "POST", "/import/allocations", importHandler.CreateAllocation, RouteSpec{ - Summary: "创建CSV设备批量分配任务", - Description: "复用设备导入任务。CSV只能包含一列设备标识,每行支持VirtualNo、IMEI或SN;整批选择分配目标代理或设置目标套餐系列。", + Summary: "创建CSV设备批量分配或回收任务", + Description: "复用设备导入任务。CSV只能包含一列设备标识,每行支持VirtualNo、IMEI或SN;整批可分配目标代理、设置目标套餐系列或回收设备。", Tags: []string{"设备管理"}, Input: new(dto.CreateDeviceBatchAllocationRequest), Output: new(dto.CreateDeviceBatchAllocationResponse), diff --git a/internal/routes/storage.go b/internal/routes/storage.go index 128aa00..dcdc4a5 100644 --- a/internal/routes/storage.go +++ b/internal/routes/storage.go @@ -135,7 +135,7 @@ await api.post('/iot-cards/import', { |---|------|-------------| | iot_import | ICCID/设备导入 (Excel) | imports/YYYY/MM/DD/uuid.xlsx | | batch_purchase | 资产套餐批量订购 (CSV) | batch-purchases/YYYY/MM/DD/uuid.csv | -| device_batch_allocation | 设备批量分配代理或套餐系列 (CSV) | device-batch-allocations/YYYY/MM/DD/uuid.csv | +| device_batch_allocation | 设备批量分配、设置套餐系列或回收 (CSV) | device-batch-allocations/YYYY/MM/DD/uuid.csv | | export | 数据导出 | exports/YYYY/MM/DD/uuid.xlsx | | attachment | 附件上传 | attachments/YYYY/MM/DD/uuid.ext | diff --git a/internal/service/device_import/audit.go b/internal/service/device_import/audit.go index ffb3186..9b755f4 100644 --- a/internal/service/device_import/audit.go +++ b/internal/service/device_import/audit.go @@ -61,13 +61,15 @@ func newDeviceBatchAllocationAuditParams(taskID uint, taskNo string, req *dto.Cr if req != nil { afterData["file_key"] = req.FileKey afterData["operation_type"] = req.OperationType - afterData["target_id"] = req.TargetID + if req.OperationType != constants.DeviceImportOperationRecall { + afterData["target_id"] = req.TargetID + } } errorCode, errorMsg := assetAuditSvc.BuildErrorInfo(err) return assetAuditSvc.BuildLogParams{ AssetID: taskID, AssetIdentifier: taskNo, OperationType: constants.AssetAuditOpDeviceBatchTaskCreate, - OperationDesc: "创建设备CSV批量分配任务", ResultStatus: resultStatus, + OperationDesc: "创建设备CSV批量操作任务", ResultStatus: resultStatus, ErrorCode: errorCode, ErrorMsg: errorMsg, AfterData: afterData, } } diff --git a/internal/service/device_import/service.go b/internal/service/device_import/service.go index 3e164cc..f15260c 100644 --- a/internal/service/device_import/service.go +++ b/internal/service/device_import/service.go @@ -96,20 +96,24 @@ func (s *Service) CreateImportTask(ctx context.Context, req *dto.ImportDeviceReq }, nil } -// CreateBatchAllocationTask 复用设备导入任务创建单列 CSV 批量分配任务。 +// CreateBatchAllocationTask 复用设备导入任务创建单列 CSV 批量操作任务。 func (s *Service) CreateBatchAllocationTask(ctx context.Context, req *dto.CreateDeviceBatchAllocationRequest) (*dto.CreateDeviceBatchAllocationResponse, error) { userID := middleware.GetUserIDFromContext(ctx) userType := middleware.GetUserTypeFromContext(ctx) if userID == 0 || (userType != constants.UserTypeSuperAdmin && userType != constants.UserTypePlatform && userType != constants.UserTypeAgent) { - appErr := errors.New(errors.CodeForbidden, "仅平台和代理后台账号可创建设备批量分配任务") + appErr := errors.New(errors.CodeForbidden, "仅平台和代理后台账号可创建设备CSV批量任务") s.logDeviceImportAudit(ctx, newDeviceBatchAllocationAuditParams(0, "", req, constants.AssetAuditResultDenied, appErr)) return nil, appErr } - if req == nil || !constants.IsDeviceImportOperation(req.OperationType) || req.OperationType == constants.DeviceImportOperationCreate || req.TargetID == 0 { - return nil, errors.New(errors.CodeInvalidParam, "设备批量分配参数不合法") + if req == nil || !constants.IsDeviceImportOperation(req.OperationType) || req.OperationType == constants.DeviceImportOperationCreate { + return nil, errors.New(errors.CodeInvalidParam, "设备CSV批量任务参数不合法") + } + if (req.OperationType == constants.DeviceImportOperationRecall && req.TargetID != 0) || + (req.OperationType != constants.DeviceImportOperationRecall && req.TargetID == 0) { + return nil, errors.New(errors.CodeInvalidParam, "设备CSV批量任务参数不合法") } if !strings.HasPrefix(req.FileKey, constants.DeviceBatchAllocationStoragePrefix+"/") || !strings.EqualFold(filepath.Ext(req.FileKey), ".csv") { - return nil, errors.New(errors.CodeInvalidParam, "设备批量分配文件必须是指定目录下的CSV文件") + return nil, errors.New(errors.CodeInvalidParam, "设备CSV批量文件必须是指定目录下的CSV文件") } taskNo := s.importTaskStore.GenerateTaskNo(ctx) @@ -121,16 +125,19 @@ func (s *Service) CreateBatchAllocationTask(ctx context.Context, req *dto.Create } operatorShopID = &shopID } - targetID := req.TargetID + var targetID *uint + if req.OperationType != constants.DeviceImportOperationRecall { + targetID = &req.TargetID + } task := &model.DeviceImportTask{ - TaskNo: taskNo, OperationType: req.OperationType, TargetID: &targetID, + TaskNo: taskNo, OperationType: req.OperationType, TargetID: targetID, OperatorType: userType, OperatorShopID: operatorShopID, Status: model.ImportTaskStatusPending, StorageKey: req.FileKey, FileName: filepath.Base(req.FileKey), CreatorName: middleware.GetUsernameFromContext(ctx), } task.Creator, task.Updater = userID, userID if err := s.importTaskStore.Create(ctx, task); err != nil { - appErr := errors.Wrap(errors.CodeDatabaseError, err, "创建设备批量分配任务失败") + appErr := errors.Wrap(errors.CodeDatabaseError, err, "创建设备CSV批量任务失败") s.logDeviceImportAudit(ctx, newDeviceBatchAllocationAuditParams(0, taskNo, req, constants.AssetAuditResultFailed, appErr)) return nil, appErr } @@ -138,13 +145,13 @@ func (s *Service) CreateBatchAllocationTask(ctx context.Context, req *dto.Create asynq.Queue(constants.QueueForTaskType(constants.TaskTypeDeviceImport)), asynq.Timeout(constants.DeviceBatchAllocationTaskTimeout)); err != nil { _ = s.importTaskStore.UpdateStatus(ctx, task.ID, model.ImportTaskStatusFailed, "任务入队失败") - appErr := errors.Wrap(errors.CodeInternalError, err, "设备批量分配任务入队失败") + appErr := errors.Wrap(errors.CodeInternalError, err, "设备CSV批量任务入队失败") s.logDeviceImportAudit(ctx, newDeviceBatchAllocationAuditParams(task.ID, taskNo, req, constants.AssetAuditResultFailed, appErr)) return nil, appErr } s.logDeviceImportAudit(ctx, newDeviceBatchAllocationAuditParams(task.ID, taskNo, req, constants.AssetAuditResultSuccess, nil)) return &dto.CreateDeviceBatchAllocationResponse{ - TaskID: task.ID, TaskNo: task.TaskNo, Message: "设备批量分配任务已创建,Worker 将异步处理CSV文件", + TaskID: task.ID, TaskNo: task.TaskNo, Message: "设备CSV批量任务已创建,Worker 将异步处理CSV文件", }, nil } diff --git a/internal/task/device_batch_allocation.go b/internal/task/device_batch_allocation.go index dfa4102..33125d6 100644 --- a/internal/task/device_batch_allocation.go +++ b/internal/task/device_batch_allocation.go @@ -12,6 +12,7 @@ import ( "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/postgres" "github.com/break/junhong_cmp_fiber/pkg/constants" "github.com/break/junhong_cmp_fiber/pkg/errors" "github.com/break/junhong_cmp_fiber/pkg/middleware" @@ -23,8 +24,8 @@ type deviceBatchAllocationRow struct { } func (h *DeviceImportHandler) handleDeviceBatchAllocation(ctx context.Context, task *model.DeviceImportTask) error { - if h.allocationExecutor == nil || task.TargetID == nil || *task.TargetID == 0 { - _ = h.importTaskStore.UpdateStatus(ctx, task.ID, model.ImportTaskStatusFailed, "设备批量分配执行器或目标未配置") + if h.allocationExecutor == nil || (task.OperationType != constants.DeviceImportOperationRecall && (task.TargetID == nil || *task.TargetID == 0)) { + _ = h.importTaskStore.UpdateStatus(ctx, task.ID, model.ImportTaskStatusFailed, "设备CSV批量执行器或目标未配置") return asynq.SkipRetry } rows, err := h.downloadDeviceBatchAllocationCSV(ctx, task) @@ -32,9 +33,14 @@ func (h *DeviceImportHandler) handleDeviceBatchAllocation(ctx context.Context, t _ = h.importTaskStore.UpdateStatus(ctx, task.ID, model.ImportTaskStatusFailed, err.Error()) return asynq.SkipRetry } + shopScope, err := h.resolveDeviceBatchShopScope(ctx, task) + if err != nil { + _ = h.importTaskStore.UpdateStatus(ctx, task.ID, model.ImportTaskStatusFailed, err.Error()) + return asynq.SkipRetry + } workerCtx := middleware.SetUserContext(ctx, &middleware.UserContextInfo{ UserID: task.Creator, UserType: task.OperatorType, Username: task.CreatorName, - ShopID: valueOrZero(task.OperatorShopID), SubordinateShopIDs: operatorDeviceShopScope(task.OperatorShopID), + ShopID: valueOrZero(task.OperatorShopID), SubordinateShopIDs: shopScope, }) result, err := h.executeDeviceBatchAllocation(workerCtx, task, rows) if err != nil { @@ -43,7 +49,7 @@ func (h *DeviceImportHandler) handleDeviceBatchAllocation(ctx context.Context, t } _ = h.importTaskStore.UpdateResult(ctx, task.ID, len(rows), result.successCount, result.skipCount, result.failCount, 0, result.skippedItems, result.failedItems, nil) if result.successCount == 0 && result.failCount > 0 { - _ = h.importTaskStore.UpdateStatus(ctx, task.ID, model.ImportTaskStatusFailed, "所有设备分配均失败") + _ = h.importTaskStore.UpdateStatus(ctx, task.ID, model.ImportTaskStatusFailed, "所有设备操作均失败") } else { _ = h.importTaskStore.UpdateStatus(ctx, task.ID, model.ImportTaskStatusCompleted, "") } @@ -52,20 +58,20 @@ func (h *DeviceImportHandler) handleDeviceBatchAllocation(ctx context.Context, t func (h *DeviceImportHandler) downloadDeviceBatchAllocationCSV(ctx context.Context, task *model.DeviceImportTask) ([]deviceBatchAllocationRow, error) { if h.storageService == nil || task.StorageKey == "" { - return nil, errors.New(errors.CodeServiceUnavailable, "设备批量分配对象存储未配置") + return nil, errors.New(errors.CodeServiceUnavailable, "设备CSV批量对象存储未配置") } path, cleanup, err := h.storageService.DownloadToTemp(ctx, task.StorageKey) if err != nil { - return nil, errors.Wrap(errors.CodeInternalError, err, "下载设备批量分配CSV失败") + return nil, errors.Wrap(errors.CodeInternalError, err, "下载设备批量CSV失败") } defer cleanup() info, err := os.Stat(path) if err != nil || info.Size() > constants.DeviceBatchAllocationMaxFileSize { - return nil, errors.New(errors.CodeInvalidParam, "设备批量分配CSV不存在或超过10MB") + return nil, errors.New(errors.CodeInvalidParam, "设备批量CSV不存在或超过10MB") } data, err := os.ReadFile(path) if err != nil { - return nil, errors.Wrap(errors.CodeInternalError, err, "读取设备批量分配CSV失败") + return nil, errors.Wrap(errors.CodeInternalError, err, "读取设备批量CSV失败") } return parseDeviceBatchAllocationCSV(data) } @@ -81,22 +87,22 @@ func parseDeviceBatchAllocationCSV(data []byte) ([]deviceBatchAllocationRow, err break } if err != nil || len(record) != 1 { - return nil, errors.New(errors.CodeInvalidParam, "设备批量分配CSV必须只有一列设备标识") + return nil, errors.New(errors.CodeInvalidParam, "设备批量CSV必须只有一列设备标识") } identifier := strings.TrimSpace(record[0]) if line == 1 && (identifier == "device_identifier" || identifier == "设备标识") { continue } if identifier == "" { - return nil, errors.New(errors.CodeInvalidParam, "设备批量分配CSV存在空设备标识") + return nil, errors.New(errors.CodeInvalidParam, "设备批量CSV存在空设备标识") } rows = append(rows, deviceBatchAllocationRow{line: line, identifier: identifier}) if len(rows) > constants.DeviceBatchAllocationMaxRows { - return nil, errors.New(errors.CodeInvalidParam, "设备批量分配CSV最多包含1000行设备") + return nil, errors.New(errors.CodeInvalidParam, "设备批量CSV最多包含1000行设备") } } if len(rows) == 0 { - return nil, errors.New(errors.CodeInvalidParam, "设备批量分配CSV没有有效数据行") + return nil, errors.New(errors.CodeInvalidParam, "设备批量CSV没有有效数据行") } return rows, nil } @@ -135,7 +141,7 @@ func (h *DeviceImportHandler) executeDeviceBatchAllocation(ctx context.Context, continue } seen[device.ID] = struct{}{} - if deviceAlreadyAtAllocationTarget(device, task.OperationType, *task.TargetID) { + if deviceAlreadyAtAllocationTarget(device, task) { result.successCount++ continue } @@ -179,17 +185,34 @@ func (h *DeviceImportHandler) applyDeviceBatchAllocation(ctx context.Context, ta for _, item := range response.FailedItems { failed[item.DeviceID] = item.Reason } + case constants.DeviceImportOperationRecall: + response, err := h.allocationExecutor.RecallDevices(ctx, &dto.RecallDevicesRequest{DeviceIDs: deviceIDs, Remark: "CSV批量回收任务 " + task.TaskNo}, task.Creator, task.OperatorShopID) + if err != nil { + return nil, err + } + for _, item := range response.FailedItems { + failed[item.DeviceID] = item.Reason + } default: - return nil, errors.New(errors.CodeInvalidParam, "设备批量分配任务类型不支持") + return nil, errors.New(errors.CodeInvalidParam, "设备CSV批量任务类型不支持") } return failed, nil } -func deviceAlreadyAtAllocationTarget(device *model.Device, operation string, targetID uint) bool { - if operation == constants.DeviceImportOperationAssignShop { - return device.ShopID != nil && *device.ShopID == targetID +func deviceAlreadyAtAllocationTarget(device *model.Device, task *model.DeviceImportTask) bool { + switch task.OperationType { + case constants.DeviceImportOperationAssignShop: + return task.TargetID != nil && device.ShopID != nil && *device.ShopID == *task.TargetID + case constants.DeviceImportOperationAssignSeries: + return task.TargetID != nil && device.SeriesID != nil && *device.SeriesID == *task.TargetID + case constants.DeviceImportOperationRecall: + if task.OperatorShopID == nil { + return device.ShopID == nil + } + return device.ShopID != nil && *device.ShopID == *task.OperatorShopID + default: + return false } - return operation == constants.DeviceImportOperationAssignSeries && device.SeriesID != nil && *device.SeriesID == targetID } func deviceAllocationResultItem(row deviceBatchAllocationRow, reason string) model.ImportResultItem { @@ -203,9 +226,16 @@ func valueOrZero(value *uint) uint { return *value } -func operatorDeviceShopScope(shopID *uint) []uint { - if shopID == nil { - return nil +func (h *DeviceImportHandler) resolveDeviceBatchShopScope(ctx context.Context, task *model.DeviceImportTask) ([]uint, error) { + if task.OperatorShopID == nil { + return nil, nil } - return []uint{*shopID} + if task.OperationType != constants.DeviceImportOperationRecall { + return []uint{*task.OperatorShopID}, nil + } + shopIDs, err := postgres.NewShopStore(h.db, h.redis).GetSubordinateShopIDs(ctx, *task.OperatorShopID) + if err != nil { + return nil, errors.Wrap(errors.CodeDatabaseError, err, "查询代理设备回收范围失败") + } + return shopIDs, nil } diff --git a/internal/task/device_import.go b/internal/task/device_import.go index a70816f..d572dc3 100644 --- a/internal/task/device_import.go +++ b/internal/task/device_import.go @@ -32,6 +32,7 @@ type DeviceImportPayload struct { type DeviceBatchAllocationExecutor interface { AllocateDevices(ctx context.Context, req *dto.AllocateDevicesRequest, operatorID uint, operatorShopID *uint) (*dto.AllocateDevicesResponse, error) BatchSetSeriesBinding(ctx context.Context, req *dto.BatchSetDeviceSeriesBindngRequest, operatorShopID *uint) (*dto.BatchSetDeviceSeriesBindngResponse, error) + RecallDevices(ctx context.Context, req *dto.RecallDevicesRequest, operatorID uint, operatorShopID *uint) (*dto.RecallDevicesResponse, error) } type DeviceImportHandler struct { diff --git a/pkg/constants/asset_audit.go b/pkg/constants/asset_audit.go index 6bf3d8f..c6980a3 100644 --- a/pkg/constants/asset_audit.go +++ b/pkg/constants/asset_audit.go @@ -50,7 +50,7 @@ const ( AssetAuditOpAssetPackageUsage = "asset_package_usage" // 资产套餐已用量更新 AssetAuditOpCardSpeedTier = "card_speed_tier" // IoT 卡固定限速档位设置 AssetAuditOpDeviceImportTaskCreate = "device_import_task_create" // 设备导入任务创建 - AssetAuditOpDeviceBatchTaskCreate = "device_batch_task_create" // 设备CSV批量分配任务创建 + AssetAuditOpDeviceBatchTaskCreate = "device_batch_task_create" // 设备 CSV 批量操作任务创建 AssetAuditOpIotCardImportTaskCreate = "iot_card_import_task_create" // 卡导入任务创建 ) diff --git a/pkg/constants/device_batch_allocation.go b/pkg/constants/device_batch_allocation.go index 324c6b1..4f144a7 100644 --- a/pkg/constants/device_batch_allocation.go +++ b/pkg/constants/device_batch_allocation.go @@ -7,12 +7,13 @@ const ( DeviceImportOperationCreate = "import" // 导入并创建设备 DeviceImportOperationAssignShop = "assign_shop" // 按单列 CSV 分配目标代理店铺 DeviceImportOperationAssignSeries = "assign_series" // 按单列 CSV 设置目标套餐系列 + DeviceImportOperationRecall = "recall" // 按单列 CSV 回收设备 ) -// 设备批量分配 CSV 约束。 +// 设备 CSV 批量操作约束。 const ( - StoragePurposeDeviceBatchAllocation = "device_batch_allocation" // 设备批量分配 CSV 上传用途 - DeviceBatchAllocationStoragePrefix = "device-batch-allocations" // 设备批量分配对象存储目录 + StoragePurposeDeviceBatchAllocation = "device_batch_allocation" // 设备 CSV 批量操作上传用途 + DeviceBatchAllocationStoragePrefix = "device-batch-allocations" // 设备 CSV 批量操作对象存储目录 DeviceBatchAllocationMaxRows = 1000 // 单个 CSV 最大设备行数 DeviceBatchAllocationMaxFileSize = int64(10 * 1024 * 1024) // 单个 CSV 最大字节数 DeviceBatchAllocationTaskTimeout = 2 * time.Hour // 单个任务最长执行时间 @@ -21,7 +22,7 @@ const ( // IsDeviceImportOperation 判断是否为设备导入任务支持的业务类型。 func IsDeviceImportOperation(operation string) bool { switch operation { - case DeviceImportOperationCreate, DeviceImportOperationAssignShop, DeviceImportOperationAssignSeries: + case DeviceImportOperationCreate, DeviceImportOperationAssignShop, DeviceImportOperationAssignSeries, DeviceImportOperationRecall: return true default: return false @@ -37,6 +38,8 @@ func GetDeviceImportOperationName(operation string) string { return "分配目标代理" case DeviceImportOperationAssignSeries: return "设置套餐系列" + case DeviceImportOperationRecall: + return "回收设备" default: return "未知" }