package repository import ( "context" "database/sql" "errors" "fmt" "strconv" "strings" "cmautobuy/admin/model" ) const ( collectTaskIDPrefix = "cj" purchaseTaskIDPrefix = "cg" ) // NextTaskID 在当前事务内分配下一条任务业务单号。 // // 采集和采购分别递增;更新序列表后再读取同一行,依赖事务行锁保证并发不重复。 // 调用方必须把分配和任务 INSERT 放在同一个事务里,失败回滚时不会消耗序号。 func NextTaskID(q Execer, taskType model.TaskType) (string, error) { prefix, ok := taskIDPrefix(taskType) if !ok { return "", fmt.Errorf("不支持为任务类型 %q 分配编号", taskType) } result, err := q.Exec(`UPDATE task_sequences SET current_value=current_value+1 WHERE task_type=?`, taskType) if err != nil { return "", fmt.Errorf("递增 %s 任务序号失败: %w", taskType, err) } affected, err := result.RowsAffected() if err != nil { return "", fmt.Errorf("读取 %s 任务序号更新结果失败: %w", taskType, err) } if affected != 1 { return "", fmt.Errorf("%s 任务序列表缺失,请先完成数据库迁移", taskType) } var number int64 if err := q.QueryRow(`SELECT current_value FROM task_sequences WHERE task_type=?`, taskType).Scan(&number); err != nil { return "", fmt.Errorf("读取 %s 任务序号失败: %w", taskType, err) } if number <= 0 { return "", fmt.Errorf("%s 任务序号不正确: %d", taskType, number) } return prefix + strconv.FormatInt(number, 10), nil } func taskIDPrefix(taskType model.TaskType) (string, bool) { switch taskType { case model.TaskCollect: return collectTaskIDPrefix, true case model.TaskPurchase: return purchaseTaskIDPrefix, true default: return "", false } } func parseTaskIDNumber(taskType model.TaskType, taskID string) (int64, bool) { prefix, ok := taskIDPrefix(taskType) if !ok || !strings.HasPrefix(taskID, prefix) || len(taskID) == len(prefix) { return 0, false } numberText := taskID[len(prefix):] if numberText[0] == '0' { return 0, false } number, err := strconv.ParseInt(numberText, 10, 64) return number, err == nil && number > 0 } // ClaimNextTask 为指定客户端领取一个任务。 // // 没有可领的任务时返回 (nil, nil) —— 调用方据此返回 204。 // // # 两种任务都能领 // // 指定给本机的 assigned_client = 我 且 status = 'assigned' // 无主的 assigned_client 为空 且 status = 'pending' // // 采集任务默认不指定客户端(浏览商品页没有副作用,哪台设备采都一样), // PDD 批量页面也允许人工指定,因此两种领取条件都必须保留。 // 采购任务可以指定也可以留空——涉及钱和账号时应当显式分配。 // // **指定给本机的优先。** 显式分配是人为决定,应当先兑现; // 无主任务谁抢都一样,可以等。 // // # 防并发 // // InnoDB 事务用 FOR UPDATE SKIP LOCKED 锁住一条候选任务。领取状态和 // task_claims 历史在同一个事务提交,避免只改了状态却没留下领取凭据。 func ClaimNextTask(db *sql.DB, clientID string, supportedTypes []string, purchaseModes ...string) (*model.Task, error) { if clientID == "" { return nil, fmt.Errorf("client_id 不能为空") } purchaseMode := string(model.TaskExecutionDryRun) if len(purchaseModes) > 0 { purchaseMode = purchaseModes[0] } if purchaseMode != string(model.TaskExecutionDryRun) && purchaseMode != string(model.TaskExecutionLive) { return nil, fmt.Errorf("purchase_mode 必须是 dry_run 或 live") } ctx := context.Background() tx, err := db.BeginTx(ctx, &sql.TxOptions{Isolation: sql.LevelReadCommitted}) if err != nil { return nil, fmt.Errorf("开始领取任务事务失败: %w", err) } defer tx.Rollback() query := `SELECT task_id FROM tasks WHERE ( (assigned_client = ? AND status = 'assigned') OR (assigned_client IS NULL AND status = 'pending') )` args := []any{clientID} // dry_run 客户端永远看不到真实采购任务。声明 live 的客户端仍可执行 // 演练任务,避免它在没有真实任务时闲置。 if purchaseMode != string(model.TaskExecutionLive) { query += ` AND execution_mode = 'dry_run'` } // 客户端只声明支持某些类型时,不要给它别的类型 if len(supportedTypes) > 0 { placeholders := strings.TrimSuffix(strings.Repeat("?,", len(supportedTypes)), ",") query += ` AND task_type IN (` + placeholders + `)` for _, t := range supportedTypes { args = append(args, t) } } // (assigned_client IS NULL) 为 0/1,0 排前面 —— 指定给本机的优先于无主的 query += ` ORDER BY (assigned_client IS NULL), priority DESC, created_at LIMIT 1 FOR UPDATE SKIP LOCKED` var taskID string if err := tx.QueryRowContext(ctx, query, args...).Scan(&taskID); errors.Is(err, sql.ErrNoRows) { return nil, nil } else if err != nil { return nil, fmt.Errorf("查询可领任务失败: %w", err) } now := model.NowISO() // 两种情况合成一条语句:对"指定给我的"那种,写 assigned_client // 是写同一个值,无副作用;对无主的,这一步就是"谁领到就标记谁"。 res, err := tx.ExecContext(ctx, ` UPDATE tasks SET status = 'claimed', assigned_client = ?, claimed_at = ?, updated_at = ? WHERE task_id = ? AND ( (status = 'assigned' AND assigned_client = ?) OR (status = 'pending' AND assigned_client IS NULL) )`, clientID, now, now, taskID, clientID) if err != nil { return nil, fmt.Errorf("领取任务 %s 失败: %w", taskID, err) } n, err := res.RowsAffected() if err != nil { return nil, fmt.Errorf("确认领取任务 %s 结果失败: %w", taskID, err) } if n != 1 { return nil, fmt.Errorf("领取任务 %s 时状态异常", taskID) } // 记一笔领取历史。提交结果时要靠它判断这台客户端有没有领过—— // 任务重派后 assigned_client 会变,只看它就查不出来了。 if err := RecordClaim(tx, taskID, clientID, now); err != nil { return nil, err } if err := tx.Commit(); err != nil { return nil, fmt.Errorf("提交领取任务事务失败: %w", err) } return GetTask(db, taskID) } // GetTask 按编号读一条任务。 func GetTask(db *sql.DB, taskID string) (*model.Task, error) { var t model.Task var assigned, claimedAt, sybID, orderNo, goodsID, skuID sql.NullString var pddGoodsID, pddOptions, resultData, errCode, errMsg, finishedAt sql.NullString var liveConfirmedBy, liveConfirmedAt, createdByUserID sql.NullString var quantity, maxPrice sql.NullInt64 err := db.QueryRow(` SELECT task_id, task_type, status, execution_mode, version, priority, assigned_client, claimed_at, syb_id, order_no, goods_id, shopee_sku_id, pdd_goods_url, pdd_goods_id, pdd_options, quantity, max_price_cent, result_data, error_code, error_message, finished_at, live_confirmed_by, live_confirmed_at, created_by_user_id, created_at, updated_at FROM tasks WHERE task_id = ?`, taskID).Scan( &t.TaskID, &t.TaskType, &t.Status, &t.ExecutionMode, &t.Version, &t.Priority, &assigned, &claimedAt, &sybID, &orderNo, &goodsID, &skuID, &t.PddGoodsURL, &pddGoodsID, &pddOptions, &quantity, &maxPrice, &resultData, &errCode, &errMsg, &finishedAt, &liveConfirmedBy, &liveConfirmedAt, &createdByUserID, &t.CreatedAt, &t.UpdatedAt) if err == sql.ErrNoRows { return nil, nil } if err != nil { return nil, fmt.Errorf("读取任务 %s 失败: %w", taskID, err) } t.AssignedClient = assigned.String t.ClaimedAt = claimedAt.String t.SybID = sybID.String t.OrderNo = orderNo.String t.GoodsID = goodsID.String t.ShopeeSKUID = skuID.String t.PddGoodsID = pddGoodsID.String t.PddOptions = pddOptions.String t.Quantity = int(quantity.Int64) t.MaxPriceCent = maxPrice.Int64 t.ResultData = resultData.String t.ErrorCode = errCode.String t.ErrorMessage = errMsg.String t.FinishedAt = finishedAt.String t.LiveConfirmedBy = liveConfirmedBy.String t.LiveConfirmedAt = liveConfirmedAt.String t.CreatedByUserID = createdByUserID.String return &t, nil } // TaskVisibleToUser 判断任务是否在采购员创建人范围内。visibleUserID 为空表示管理员。 func TaskVisibleToUser(q Execer, taskID, visibleUserID string) (bool, error) { query := `SELECT COUNT(*) FROM tasks WHERE task_id = ?` args := []any{taskID} if visibleUserID != "" { query += ` AND created_by_user_id = ?` args = append(args, visibleUserID) } var count int if err := q.QueryRow(query, args...).Scan(&count); err != nil { return false, fmt.Errorf("检查任务可见范围失败: %w", err) } return count == 1, nil } // TaskDisplayContext 是任务详情按当前关联数据展示的名称。 // 名称可能被修改或对应记录被删除,因此 tasks 仍只保存稳定 ID。 type TaskDisplayContext struct { CreatorUsername string ClientName string PddShopName string } // GetTaskDisplayContext 返回详情需要的当前创建人、客户端名称和 PDD 店铺。 // PDD 商品软删除后不再作为当前商品档案展示,店铺名自然为空。 func GetTaskDisplayContext(q Execer, taskID string) (TaskDisplayContext, error) { var result TaskDisplayContext var creatorUsername, clientName, pddShopName sql.NullString err := q.QueryRow(` SELECT u.username, c.name, p.shop_name FROM tasks t LEFT JOIN users u ON u.user_id = t.created_by_user_id LEFT JOIN clients c ON c.client_id = t.assigned_client LEFT JOIN pdd_products p ON p.goods_id = t.pdd_goods_id AND p.deleted_at IS NULL WHERE t.task_id = ?`, taskID).Scan(&creatorUsername, &clientName, &pddShopName) if err != nil { if err == sql.ErrNoRows { return result, nil } return result, fmt.Errorf("读取任务展示信息失败: %w", err) } result.CreatorUsername = creatorUsername.String result.ClientName = clientName.String result.PddShopName = pddShopName.String return result, nil } // RecordClaim 记一笔"某客户端领过某任务"。 // // 同一台客户端重复领同一个任务时只更新时间,不报错。 func RecordClaim(q Execer, taskID, clientID, claimedAt string) error { _, err := q.Exec(` INSERT INTO task_claims (task_id, client_id, claimed_at) VALUES (?, ?, ?) ON DUPLICATE KEY UPDATE claimed_at = VALUES(claimed_at)`, taskID, clientID, claimedAt) if err != nil { return fmt.Errorf("记录领取历史失败 task=%s client=%s: %w", taskID, clientID, err) } return nil } // HasEverClaimed 判断这台客户端**曾经**领过这个任务。 // // 用于提交结果时的权限判断:只有从没领过的才拒绝(403)。 // 任务已取消、已重派给别人,都**不影响**这个判断—— // 契约要求那些情况仍然要接受结果,见 docs/admin/04-client-api.md §4.1。 func HasEverClaimed(q Execer, taskID, clientID string) (bool, error) { var one int err := q.QueryRow( `SELECT 1 FROM task_claims WHERE task_id = ? AND client_id = ?`, taskID, clientID).Scan(&one) if err == sql.ErrNoRows { return false, nil } if err != nil { return false, fmt.Errorf("查询领取历史失败: %w", err) } return true, nil } // TaskInfo 是提交结果时需要知道的任务基本信息。 type TaskInfo struct { TaskType model.TaskType // GoodsID 是关联的**蝦皮**商品编号。 GoodsID string // PddGoodsID 是要采集/购买的**拼多多**商品编号。 // 采集结果落到 pdd_products 时用的是它,不是 GoodsID —— 被采集的是 PDD 商品。 PddGoodsID string } // GetTaskInfo 查任务的类型和关联商品。任务不存在时返回 (nil, nil)。 func GetTaskInfo(q Execer, taskID string) (*TaskInfo, error) { var info TaskInfo var gid, pddGID sql.NullString err := q.QueryRow( `SELECT task_type, goods_id, pdd_goods_id FROM tasks WHERE task_id = ?`, taskID).Scan(&info.TaskType, &gid, &pddGID) if err == sql.ErrNoRows { return nil, nil } if err != nil { return nil, fmt.Errorf("查询任务 %s 失败: %w", taskID, err) } info.GoodsID = gid.String info.PddGoodsID = pddGID.String return &info, nil } // MarkTaskSucceeded 记录成功结果。 // // `[必须]` **不检查任务当前状态**。任务已取消、已重派给别的客户端, // 都照样接受——客户端可能真的已经下单了,这些数据必须能留痕。 // 理由见 docs/admin/04-client-api.md §4.1。 func MarkTaskSucceeded(q Execer, taskID, resultData string) error { now := model.NowISO() _, err := q.Exec(` UPDATE tasks SET status = 'succeeded', result_data = ?, finished_at = ?, updated_at = ? WHERE task_id = ?`, resultData, now, now, taskID) if err != nil { return fmt.Errorf("标记任务 %s 成功失败: %w", taskID, err) } return nil } // MarkTaskFailure 记录失败/需人工处理的结果。 // // status 由客户端报告,映射规则见 docs/admin/04-client-api.md §5: // // retry_wait -> assigned(放回去等它再来领) // manual_review -> manual_review // failed -> failed // cancelled -> cancelled // // 同样**不检查任务当前状态**,理由同 MarkTaskSucceeded。 func MarkTaskFailure(q Execer, taskID string, newStatus model.TaskStatus, errCode, errMsg string) error { now := model.NowISO() // 放回待领取的话不算结束,finished_at 保持为空 finishedAt := any(now) if newStatus == model.TaskAssigned { finishedAt = nil } _, err := q.Exec(` UPDATE tasks SET status = ?, error_code = ?, error_message = ?, finished_at = ?, updated_at = ? WHERE task_id = ?`, newStatus, errCode, errMsg, finishedAt, now, taskID) if err != nil { return fmt.Errorf("标记任务 %s 失败状态出错: %w", taskID, err) } return nil } // ---------- 采集采购列表(#19) ---------- // TaskFilter 是「采集采购」列表页支持的筛选条件,三项都可以为空。 type TaskFilter struct { Type model.TaskType Status model.TaskStatus Keyword string // 同时匹配任务编号、订单号、PDD 商品 ID // VisibleUserID 非空时强制只返回该创建人的任务,供采购员权限隔离。 VisibleUserID string // CreatorUserID / CreatorHistory 只供管理员的创建人筛选使用。 CreatorUserID string CreatorHistory bool } // TaskListRow 是列表一行要用到的原始字段,还没翻成界面文字—— // 那是 service 层的事(尤其是「目标」列的拼接,见 #19)。 type TaskListRow struct { TaskID string TaskType model.TaskType Status model.TaskStatus ExecutionMode model.TaskExecutionMode AssignedClient string // 空表示无主任务 CreatedByUserID string CreatedByUsername string ClientName string OrderNo string PddGoodsID string PddOptions string Quantity int MaxPriceCent int64 UpdatedAt string // PddTitle 是 join pdd_products 拿到的标题,**可能为空**: // 没有对应商品、或商品已被软删除时都是空。 // `[必须]` service 层退回显示 PddGoodsID,不能显示空,见 #19。 PddTitle string PddShopName string } // taskFilterClause 把三项筛选条件拼成 WHERE 子句,供 ListTasks 和 // CountTasksByStatus 共用——两处筛选逻辑必须完全一致, // 否则底部统计会跟表格对不上(见 #19 的验收要求)。 func taskFilterClause(filter TaskFilter) (string, []any) { var clauses []string var args []any if filter.Type != "" { clauses = append(clauses, "t.task_type = ?") args = append(args, string(filter.Type)) } if filter.Status != "" { clauses = append(clauses, "t.status = ?") args = append(args, string(filter.Status)) } if kw := strings.TrimSpace(filter.Keyword); kw != "" { pattern := "%" + escapeLike(kw) + "%" clauses = append(clauses, `(t.task_id LIKE ? ESCAPE '!' OR t.order_no LIKE ? ESCAPE '!' OR t.pdd_goods_id LIKE ? ESCAPE '!' OR EXISTS ( SELECT 1 FROM task_syb_sources tss JOIN syb_orders so ON so.syb_id=tss.syb_id WHERE tss.task_id=t.task_id AND (tss.syb_id LIKE ? ESCAPE '!' OR so.order_no LIKE ? ESCAPE '!') ))`) args = append(args, pattern, pattern, pattern, pattern, pattern) } if filter.VisibleUserID != "" { clauses = append(clauses, "t.created_by_user_id = ?") args = append(args, filter.VisibleUserID) } else if filter.CreatorHistory { clauses = append(clauses, "t.created_by_user_id IS NULL") } else if filter.CreatorUserID != "" { clauses = append(clauses, "t.created_by_user_id = ?") args = append(args, filter.CreatorUserID) } if len(clauses) == 0 { return "", args } return " WHERE " + strings.Join(clauses, " AND "), args } // ListTasks 查采集采购列表,按更新时间倒序。 // // `[必须]` 目标列所需字段全在 tasks 表上,不需要 join。 // `[建议]` 这里额外 LEFT JOIN 了 pdd_products 取标题,纯粹是为了让采集任务的 // 「目标」列更好读;join 条件里的 `deleted_at IS NULL` 让"join 不到"和 // "商品已软删除"这两种情况都自然落到 PddTitle 为空,service 层退回显示 // PddGoodsID 即可,不用在 service 里再判断一次已删除。 func ListTasks(q Execer, filter TaskFilter, limit, offset int) ([]TaskListRow, error) { where, args := taskFilterClause(filter) sqlText := ` SELECT t.task_id, t.task_type, t.status, t.execution_mode, t.assigned_client, t.order_no, t.pdd_goods_id, t.pdd_options, t.quantity, t.max_price_cent, t.updated_at, p.title, p.shop_name, t.created_by_user_id, u.username, c.name FROM tasks t LEFT JOIN pdd_products p ON p.goods_id = t.pdd_goods_id AND p.deleted_at IS NULL LEFT JOIN users u ON u.user_id = t.created_by_user_id LEFT JOIN clients c ON c.client_id = t.assigned_client` + where + ` ORDER BY t.updated_at DESC, t.task_id DESC LIMIT ? OFFSET ?` args = append(args, limit, offset) rows, err := q.Query(sqlText, args...) if err != nil { return nil, fmt.Errorf("查询任务列表失败: %w", err) } defer rows.Close() list := make([]TaskListRow, 0, 16) for rows.Next() { var r TaskListRow var assigned, orderNo, pddGoodsID, pddOptions, title, shopName sql.NullString var creatorID, creatorName, clientName sql.NullString var quantity, maxPrice sql.NullInt64 if err := rows.Scan( &r.TaskID, &r.TaskType, &r.Status, &r.ExecutionMode, &assigned, &orderNo, &pddGoodsID, &pddOptions, &quantity, &maxPrice, &r.UpdatedAt, &title, &shopName, &creatorID, &creatorName, &clientName, ); err != nil { return nil, fmt.Errorf("读取任务列表失败: %w", err) } r.AssignedClient = assigned.String r.OrderNo = orderNo.String r.PddGoodsID = pddGoodsID.String r.PddOptions = pddOptions.String r.Quantity = int(quantity.Int64) r.MaxPriceCent = maxPrice.Int64 r.PddTitle = title.String r.PddShopName = shopName.String r.CreatedByUserID = creatorID.String r.CreatedByUsername = creatorName.String r.ClientName = clientName.String list = append(list, r) } return list, rows.Err() } // CountTasksByStatus 按当前筛选统计各状态的任务数。 // // `[必须]` **用和 ListTasks 完全相同的筛选条件**——统计要跟随当前筛选, // 筛了「采集」就只统计采集任务,否则底部数字和表格对不上,见 #19。 func CountTasksByStatus(q Execer, filter TaskFilter) (map[model.TaskStatus]int, error) { where, args := taskFilterClause(filter) sqlText := `SELECT t.status, COUNT(*) FROM tasks t` + where + ` GROUP BY t.status` rows, err := q.Query(sqlText, args...) if err != nil { return nil, fmt.Errorf("统计任务状态失败: %w", err) } defer rows.Close() counts := map[model.TaskStatus]int{} for rows.Next() { var status string var n int if err := rows.Scan(&status, &n); err != nil { return nil, err } counts[model.TaskStatus(status)] = n } return counts, rows.Err() } // DeleteTasks 按任务编号批量删除,返回实际删掉的条数。 // // `[必须]` 硬删,不是软删除:tasks 表本来就没有 deleted_at 列, // 加一列属于数据库结构变更,不在本工单范围内(见 #19「不做」清单)。 // 遗留下来的 task_claims 记录不删——它没有外键约束,留着不影响任何查询, // 之后想清理是可以独立做的小事,不值得为它扩大这次改动的范围。 func DeleteTasks(q Execer, taskIDs []string) (int64, error) { if len(taskIDs) == 0 { return 0, nil } placeholders := strings.TrimSuffix(strings.Repeat("?,", len(taskIDs)), ",") args := make([]any, len(taskIDs)) for i, id := range taskIDs { args[i] = id } res, err := q.Exec(`DELETE FROM tasks WHERE task_id IN (`+placeholders+`)`, args...) if err != nil { return 0, fmt.Errorf("批量删除任务失败: %w", err) } return res.RowsAffected() } // DeleteTasksInScope 原子删除指定权限范围内的任务。只要任一编号不存在或不在 // 当前采购员范围内,受影响行数就不足,调用方必须回滚整个事务。 func DeleteTasksInScope(q Execer, taskIDs []string, visibleUserID string) (int64, error) { if len(taskIDs) == 0 { return 0, nil } placeholders := strings.TrimSuffix(strings.Repeat("?,", len(taskIDs)), ",") args := make([]any, 0, len(taskIDs)+1) for _, id := range taskIDs { args = append(args, id) } where := `task_id IN (` + placeholders + `)` if visibleUserID != "" { where += ` AND created_by_user_id = ?` args = append(args, visibleUserID) } res, err := q.Exec(`DELETE FROM tasks WHERE `+where, args...) if err != nil { return 0, fmt.Errorf("批量删除任务失败: %w", err) } return res.RowsAffected() } // ListCollectTaskGoodsIDsInScope 在删除前找出受影响的 PDD 商品。 // 调用方必须和删除使用同一事务、同一可见范围,避免越权数据影响状态回收。 func ListCollectTaskGoodsIDsInScope(q Execer, taskIDs []string, visibleUserID string) ([]string, error) { if len(taskIDs) == 0 { return nil, nil } placeholders := strings.TrimSuffix(strings.Repeat("?,", len(taskIDs)), ",") args := make([]any, 0, len(taskIDs)+1) for _, id := range taskIDs { args = append(args, id) } where := `task_id IN (` + placeholders + `) AND task_type='collect' AND pdd_goods_id IS NOT NULL` if visibleUserID != "" { where += ` AND created_by_user_id = ?` args = append(args, visibleUserID) } rows, err := q.Query(`SELECT DISTINCT pdd_goods_id FROM tasks WHERE `+where, args...) if err != nil { return nil, fmt.Errorf("查询待删除采集任务的 PDD 商品失败: %w", err) } defer rows.Close() var goodsIDs []string for rows.Next() { var goodsID string if err := rows.Scan(&goodsID); err != nil { return nil, fmt.Errorf("读取待删除采集任务的 PDD 商品失败: %w", err) } goodsIDs = append(goodsIDs, goodsID) } return goodsIDs, rows.Err() } // ResetCollectingIfNoActiveTask 把没有有效采集任务的孤儿状态回收到待采集。 // 条件判断和更新在一条 SQL 内完成,避免并发建任务时误覆盖 collecting。 func ResetCollectingIfNoActiveTask(q Execer, goodsID string) (bool, error) { result, err := q.Exec(`UPDATE pdd_products AS pp SET collect_status='pending', collect_msg=NULL, updated_at=? WHERE pp.goods_id=? AND pp.collect_status='collecting' AND NOT EXISTS ( SELECT 1 FROM tasks t WHERE t.task_type='collect' AND t.pdd_goods_id=pp.goods_id AND t.status IN ('pending','assigned','claimed') )`, model.NowISO(), goodsID) if err != nil { return false, fmt.Errorf("回收商品 %s 的孤立采集中状态失败: %w", goodsID, err) } n, err := result.RowsAffected() if err != nil { return false, fmt.Errorf("确认商品 %s 的采集状态回收结果失败: %w", goodsID, err) } return n == 1, nil } // InsertCollectTask 建一条不指定客户端的采集任务。 // 保留这个入口给蝦皮、顺运宝等现有流程使用,避免它们被 PDD 页的新选项影响。 func InsertCollectTask(q Execer, taskID, goodsID, goodsURL string) error { return InsertCollectTaskForClientAndUser(q, taskID, goodsID, goodsURL, "", "") } // InsertCollectTaskForClient 建一条采集任务。 // assignedClient 为空时任务无主待领;有值时只等待指定客户端领取。 // goodsURL 必填 —— Client 契约里 pdd_goods_url 是 NOT NULL。 func InsertCollectTaskForClient(q Execer, taskID, goodsID, goodsURL, assignedClient string) error { return InsertCollectTaskForClientAndUser(q, taskID, goodsID, goodsURL, assignedClient, "") } // InsertCollectTaskForClientAndUser 建一条带创建人审计的采集任务。 func InsertCollectTaskForClientAndUser(q Execer, taskID, goodsID, goodsURL, assignedClient, createdByUserID string) error { if goodsID == "" || goodsURL == "" { return fmt.Errorf("采集任务的商品 ID 和链接都不能为空") } assignedClient = strings.TrimSpace(assignedClient) status := model.TaskPending assigned := sql.NullString{} if assignedClient != "" { status = model.TaskAssigned assigned = sql.NullString{String: assignedClient, Valid: true} } now := model.NowISO() creator := sql.NullString{String: strings.TrimSpace(createdByUserID), Valid: strings.TrimSpace(createdByUserID) != ""} _, err := q.Exec(` INSERT INTO tasks (task_id, task_type, status, assigned_client, pdd_goods_url, pdd_goods_id, created_by_user_id, created_at, updated_at) VALUES (?, 'collect', ?, ?, ?, ?, ?, ?, ?)`, taskID, status, assigned, goodsURL, goodsID, creator, now, now) if err != nil { return fmt.Errorf("创建商品 %s 的采集任务失败: %w", goodsID, err) } return nil } // InsertCollectTaskSybSources 保存一条采集任务对应的全部顺运宝明细来源。 // 同一个 PDD 商品可能由多条明细共同发起,来源不能压成 tasks 上的单列。 func InsertCollectTaskSybSources(q Execer, taskID string, sybIDs []string) error { taskID = strings.TrimSpace(taskID) if taskID == "" { return fmt.Errorf("采集任务编号不能为空") } seen := make(map[string]struct{}, len(sybIDs)) for _, raw := range sybIDs { sybID := strings.TrimSpace(raw) if sybID == "" { return fmt.Errorf("顺运宝明细编号不能为空") } if _, ok := seen[sybID]; ok { continue } seen[sybID] = struct{}{} if _, err := q.Exec(`INSERT INTO task_syb_sources(task_id,syb_id,created_at) VALUES(?,?,?) ON DUPLICATE KEY UPDATE created_at=created_at`, taskID, sybID, model.NowISO()); err != nil { return fmt.Errorf("保存采集任务 %s 的顺运宝来源 %s 失败: %w", taskID, sybID, err) } } return nil } // HasActivePurchaseTask 判断顺运宝明细是否已有尚未结束的采购任务。 func HasActivePurchaseTask(q Execer, sybID string) (bool, error) { var count int err := q.QueryRow(` SELECT COUNT(*) FROM tasks WHERE task_type = 'purchase' AND syb_id = ? AND status IN ('pending', 'assigned', 'claimed')`, sybID).Scan(&count) if err != nil { return false, fmt.Errorf("检查顺运宝明细 %s 的进行中采购任务失败: %w", sybID, err) } return count > 0, nil } // InsertPurchaseTask 插入一条已分配、等待指定客户端领取的采购任务。 func InsertPurchaseTask(q Execer, task model.Task) error { if task.TaskID == "" || task.AssignedClient == "" || task.SybID == "" || task.PddGoodsURL == "" || task.PddGoodsID == "" || task.PddOptions == "" || task.Quantity <= 0 || task.MaxPriceCent <= 0 { return fmt.Errorf("采购任务缺少客户端、商品、规格、数量或人民币价格上限") } if task.ExecutionMode == "" { task.ExecutionMode = model.TaskExecutionDryRun } if task.ExecutionMode != model.TaskExecutionDryRun && task.ExecutionMode != model.TaskExecutionLive { return fmt.Errorf("采购任务执行模式无效") } if task.ExecutionMode == model.TaskExecutionLive && (strings.TrimSpace(task.LiveConfirmedBy) == "" || strings.TrimSpace(task.LiveConfirmedAt) == "") { return fmt.Errorf("真实采购任务缺少创建确认审计信息") } if task.ExecutionMode == model.TaskExecutionDryRun { task.LiveConfirmedBy = "" task.LiveConfirmedAt = "" } now := model.NowISO() liveConfirmedBy := sql.NullString{String: task.LiveConfirmedBy, Valid: task.LiveConfirmedBy != ""} liveConfirmedAt := sql.NullString{String: task.LiveConfirmedAt, Valid: task.LiveConfirmedAt != ""} _, err := q.Exec(` INSERT INTO tasks (task_id, task_type, status, execution_mode, assigned_client, syb_id, order_no, goods_id, shopee_sku_id, pdd_goods_url, pdd_goods_id, pdd_options, quantity, max_price_cent, live_confirmed_by, live_confirmed_at, created_by_user_id, created_at, updated_at) VALUES (?, 'purchase', 'assigned', ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)`, task.TaskID, task.ExecutionMode, task.AssignedClient, task.SybID, task.OrderNo, task.GoodsID, task.ShopeeSKUID, task.PddGoodsURL, task.PddGoodsID, task.PddOptions, task.Quantity, task.MaxPriceCent, liveConfirmedBy, liveConfirmedAt, sql.NullString{String: task.CreatedByUserID, Valid: strings.TrimSpace(task.CreatedByUserID) != ""}, now, now) if err != nil { return fmt.Errorf("创建顺运宝明细 %s 的采购任务失败: %w", task.SybID, err) } return nil }