2026-08-06 16:56:18 +08:00
package repository
import (
2026-08-10 02:21:31 +08:00
"context"
2026-08-06 16:56:18 +08:00
"database/sql"
2026-08-10 02:21:31 +08:00
"errors"
2026-08-06 16:56:18 +08:00
"fmt"
"strings"
"cmautobuy/admin/model"
)
// ClaimNextTask 为指定客户端领取一个任务。
//
// 没有可领的任务时返回 (nil, nil) —— 调用方据此返回 204。
//
2026-08-07 10:53:32 +08:00
// # 两种任务都能领
//
// 指定给本机的 assigned_client = 我 且 status = 'assigned'
// 无主的 assigned_client 为空 且 status = 'pending'
//
2026-08-10 00:49:51 +08:00
// 采集任务默认不指定客户端(浏览商品页没有副作用,哪台设备采都一样),
// PDD 批量页面也允许人工指定,因此两种领取条件都必须保留。
2026-08-07 10:53:32 +08:00
// 采购任务可以指定也可以留空——涉及钱和账号时应当显式分配。
//
// **指定给本机的优先。** 显式分配是人为决定,应当先兑现;
// 无主任务谁抢都一样,可以等。
//
// # 防并发
//
2026-08-10 02:21:31 +08:00
// InnoDB 事务用 FOR UPDATE SKIP LOCKED 锁住一条候选任务。领取状态和
// task_claims 历史在同一个事务提交,避免只改了状态却没留下领取凭据。
2026-08-10 15:11:38 +08:00
func ClaimNextTask ( db * sql . DB , clientID string , supportedTypes [] string , purchaseModes ... string ) ( * model . Task , error ) {
2026-08-06 16:56:18 +08:00
if clientID == "" {
return nil , fmt . Errorf ( "client_id 不能为空" )
}
2026-08-10 15:11:38 +08:00
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" )
}
2026-08-10 02:21:31 +08:00
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 ()
2026-08-06 16:56:18 +08:00
query := `SELECT task_id FROM tasks
2026-08-07 10:53:32 +08:00
WHERE ( (assigned_client = ? AND status = 'assigned')
OR (assigned_client IS NULL AND status = 'pending') )`
2026-08-06 16:56:18 +08:00
args := [] any { clientID }
2026-08-10 15:11:38 +08:00
// dry_run 客户端永远看不到真实采购任务。声明 live 的客户端仍可执行
// 演练任务,避免它在没有真实任务时闲置。
if purchaseMode != string ( model . TaskExecutionLive ) {
query += ` AND execution_mode = 'dry_run'`
}
2026-08-06 16:56:18 +08:00
// 客户端只声明支持某些类型时,不要给它别的类型
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 )
}
}
2026-08-07 10:53:32 +08:00
// (assigned_client IS NULL) 为 0/1,0 排前面 —— 指定给本机的优先于无主的
2026-08-10 02:21:31 +08:00
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 {
2026-08-06 16:56:18 +08:00
return nil , fmt . Errorf ( "查询可领任务失败: %w" , err )
}
now := model . NowISO ()
2026-08-10 02:21:31 +08:00
// 两种情况合成一条语句:对"指定给我的"那种,写 assigned_client
// 是写同一个值,无副作用;对无主的,这一步就是"谁领到就标记谁"。
res , err := tx . ExecContext ( ctx , `
2026-08-07 10:53:32 +08:00
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) )` ,
2026-08-10 02:21:31 +08:00
clientID , now , now , taskID , clientID )
if err != nil {
return nil , fmt . Errorf ( "领取任务 %s 失败: %w" , taskID , err )
2026-08-06 16:56:18 +08:00
}
2026-08-10 02:21:31 +08:00
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 )
2026-08-06 16:56:18 +08:00
}
// 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
2026-08-10 23:10:28 +08:00
var liveConfirmedBy , liveConfirmedAt , createdByUserID sql . NullString
2026-08-06 16:56:18 +08:00
var quantity , maxPrice sql . NullInt64
err := db . QueryRow ( `
2026-08-10 15:11:38 +08:00
SELECT task_id, task_type, status, execution_mode, version, priority,
2026-08-06 16:56:18 +08:00
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,
2026-08-10 23:10:28 +08:00
live_confirmed_by, live_confirmed_at, created_by_user_id,
2026-08-06 16:56:18 +08:00
created_at, updated_at
FROM tasks WHERE task_id = ?` , taskID ). Scan (
2026-08-10 15:11:38 +08:00
& t . TaskID , & t . TaskType , & t . Status , & t . ExecutionMode , & t . Version , & t . Priority ,
2026-08-06 16:56:18 +08:00
& assigned , & claimedAt ,
& sybID , & orderNo , & goodsID , & skuID ,
& t . PddGoodsURL , & pddGoodsID , & pddOptions ,
& quantity , & maxPrice ,
& resultData , & errCode , & errMsg , & finishedAt ,
2026-08-10 23:10:28 +08:00
& liveConfirmedBy , & liveConfirmedAt , & createdByUserID ,
2026-08-06 16:56:18 +08:00
& 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
2026-08-10 15:11:38 +08:00
t . LiveConfirmedBy = liveConfirmedBy . String
t . LiveConfirmedAt = liveConfirmedAt . String
2026-08-10 23:10:28 +08:00
t . CreatedByUserID = createdByUserID . String
2026-08-06 16:56:18 +08:00
return & t , nil
}
2026-08-06 17:04:49 +08:00
2026-08-10 23:10:28 +08:00
// 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
}
// GetTaskCreatorUsername 返回任务创建人的用户名;历史任务或账号不存在时为空。
func GetTaskCreatorUsername ( q Execer , taskID string ) ( string , error ) {
var username sql . NullString
if err := q . QueryRow ( `SELECT u.username FROM tasks t LEFT JOIN users u ON u.user_id=t.created_by_user_id WHERE t.task_id=?` , taskID ). Scan ( & username ); err != nil {
if err == sql . ErrNoRows {
return "" , nil
}
return "" , fmt . Errorf ( "读取任务创建人失败: %w" , err )
}
return username . String , nil
}
2026-08-06 17:04:49 +08:00
// RecordClaim 记一笔"某客户端领过某任务"。
//
// 同一台客户端重复领同一个任务时只更新时间,不报错。
func RecordClaim ( q Execer , taskID , clientID , claimedAt string ) error {
_ , err := q . Exec ( `
INSERT INTO task_claims (task_id, client_id, claimed_at)
VALUES (?, ?, ?)
2026-08-10 02:21:31 +08:00
ON DUPLICATE KEY UPDATE claimed_at = VALUES(claimed_at)` ,
2026-08-06 17:04:49 +08:00
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
}
2026-08-07 10:27:39 +08:00
// 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
2026-08-06 17:04:49 +08:00
}
2026-08-07 10:27:39 +08:00
if err != nil {
return nil , fmt . Errorf ( "查询任务 %s 失败: %w" , taskID , err )
2026-08-06 17:04:49 +08:00
}
2026-08-07 10:27:39 +08:00
info . GoodsID = gid . String
info . PddGoodsID = pddGID . String
return & info , nil
2026-08-06 17:04:49 +08:00
}
// 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
}
2026-08-07 11:27:06 +08:00
2026-08-07 15:16:05 +08:00
// ---------- 采集采购列表(#19) ----------
// TaskFilter 是「采集采购」列表页支持的筛选条件,三项都可以为空。
type TaskFilter struct {
Type model . TaskType
Status model . TaskStatus
Keyword string // 同时匹配任务编号、订单号、PDD 商品 ID
2026-08-10 23:10:28 +08:00
// VisibleUserID 非空时强制只返回该创建人的任务,供采购员权限隔离。
VisibleUserID string
// CreatorUserID / CreatorHistory 只供管理员的创建人筛选使用。
CreatorUserID string
CreatorHistory bool
2026-08-07 15:16:05 +08:00
}
// TaskListRow 是列表一行要用到的原始字段,还没翻成界面文字——
// 那是 service 层的事(尤其是「目标」列的拼接,见 #19)。
type TaskListRow struct {
2026-08-10 23:10:28 +08:00
TaskID string
TaskType model . TaskType
Status model . TaskStatus
ExecutionMode model . TaskExecutionMode
AssignedClient string // 空表示无主任务
CreatedByUserID string
CreatedByUsername string
2026-08-07 15:16:05 +08:00
OrderNo string
PddGoodsID string
PddOptions string
Quantity int
MaxPriceCent int64
UpdatedAt string
// PddTitle 是 join pdd_products 拿到的标题,**可能为空**:
// 没有对应商品、或商品已被软删除时都是空。
// `[必须]` service 层退回显示 PddGoodsID,不能显示空,见 #19。
PddTitle 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 ,
2026-08-10 02:21:31 +08:00
"(t.task_id LIKE ? ESCAPE '!' OR t.order_no LIKE ? ESCAPE '!' OR t.pdd_goods_id LIKE ? ESCAPE '!')" )
2026-08-07 15:16:05 +08:00
args = append ( args , pattern , pattern , pattern )
}
2026-08-10 23:10:28 +08:00
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 )
}
2026-08-07 15:16:05 +08:00
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 里再判断一次已删除。
2026-08-09 21:51:20 +08:00
func ListTasks ( q Execer , filter TaskFilter , limit , offset int ) ([] TaskListRow , error ) {
2026-08-07 15:16:05 +08:00
where , args := taskFilterClause ( filter )
sqlText := `
2026-08-10 15:11:38 +08:00
SELECT t.task_id, t.task_type, t.status, t.execution_mode, t.assigned_client,
2026-08-07 15:16:05 +08:00
t.order_no, t.pdd_goods_id, t.pdd_options, t.quantity, t.max_price_cent,
2026-08-10 23:10:28 +08:00
t.updated_at, p.title, t.created_by_user_id, u.username
2026-08-07 15:16:05 +08:00
FROM tasks t
LEFT JOIN pdd_products p
2026-08-10 23:10:28 +08:00
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` +
2026-08-09 21:51:20 +08:00
where + ` ORDER BY t.updated_at DESC, t.task_id DESC LIMIT ? OFFSET ?`
args = append ( args , limit , offset )
2026-08-07 15:16:05 +08:00
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
2026-08-10 23:10:28 +08:00
var assigned , orderNo , pddGoodsID , pddOptions , title , creatorID , creatorName sql . NullString
2026-08-07 15:16:05 +08:00
var quantity , maxPrice sql . NullInt64
if err := rows . Scan (
2026-08-10 15:11:38 +08:00
& r . TaskID , & r . TaskType , & r . Status , & r . ExecutionMode , & assigned ,
2026-08-07 15:16:05 +08:00
& orderNo , & pddGoodsID , & pddOptions , & quantity , & maxPrice ,
2026-08-10 23:10:28 +08:00
& r . UpdatedAt , & title , & creatorID , & creatorName ,
2026-08-07 15:16:05 +08:00
); 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
2026-08-10 23:10:28 +08:00
r . CreatedByUserID = creatorID . String
r . CreatedByUsername = creatorName . String
2026-08-07 15:16:05 +08:00
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 ()
}
2026-08-10 23:10:28 +08:00
// 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 ()
}
2026-08-10 00:49:51 +08:00
// InsertCollectTask 建一条不指定客户端的采集任务。
// 保留这个入口给蝦皮、顺运宝等现有流程使用,避免它们被 PDD 页的新选项影响。
2026-08-07 11:27:06 +08:00
func InsertCollectTask ( q Execer , taskID , goodsID , goodsURL string ) error {
2026-08-10 23:10:28 +08:00
return InsertCollectTaskForClientAndUser ( q , taskID , goodsID , goodsURL , "" , "" )
2026-08-10 00:49:51 +08:00
}
// InsertCollectTaskForClient 建一条采集任务。
// assignedClient 为空时任务无主待领;有值时只等待指定客户端领取。
// goodsURL 必填 —— Client 契约里 pdd_goods_url 是 NOT NULL。
func InsertCollectTaskForClient ( q Execer , taskID , goodsID , goodsURL , assignedClient string ) error {
2026-08-10 23:10:28 +08:00
return InsertCollectTaskForClientAndUser ( q , taskID , goodsID , goodsURL , assignedClient , "" )
}
// InsertCollectTaskForClientAndUser 建一条带创建人审计的采集任务。
func InsertCollectTaskForClientAndUser ( q Execer , taskID , goodsID , goodsURL , assignedClient , createdByUserID string ) error {
2026-08-07 11:27:06 +08:00
if goodsID == "" || goodsURL == "" {
return fmt . Errorf ( "采集任务的商品 ID 和链接都不能为空" )
}
2026-08-10 00:49:51 +08:00
assignedClient = strings . TrimSpace ( assignedClient )
status := model . TaskPending
assigned := sql . NullString {}
if assignedClient != "" {
status = model . TaskAssigned
assigned = sql . NullString { String : assignedClient , Valid : true }
}
2026-08-07 11:27:06 +08:00
now := model . NowISO ()
2026-08-10 23:10:28 +08:00
creator := sql . NullString { String : strings . TrimSpace ( createdByUserID ), Valid : strings . TrimSpace ( createdByUserID ) != "" }
2026-08-07 11:27:06 +08:00
_ , err := q . Exec ( `
INSERT INTO tasks (task_id, task_type, status, assigned_client,
2026-08-10 23:10:28 +08:00
pdd_goods_url, pdd_goods_id, created_by_user_id, created_at, updated_at)
VALUES (?, 'collect', ?, ?, ?, ?, ?, ?, ?)` ,
taskID , status , assigned , goodsURL , goodsID , creator , now , now )
2026-08-07 11:27:06 +08:00
if err != nil {
return fmt . Errorf ( "创建商品 %s 的采集任务失败: %w" , goodsID , err )
}
return nil
}
2026-08-09 22:55:14 +08:00
// 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 ( "采购任务缺少客户端、商品、规格、数量或人民币价格上限" )
}
2026-08-10 15:11:38 +08:00
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 = ""
}
2026-08-09 22:55:14 +08:00
now := model . NowISO ()
2026-08-10 15:11:38 +08:00
liveConfirmedBy := sql . NullString { String : task . LiveConfirmedBy , Valid : task . LiveConfirmedBy != "" }
liveConfirmedAt := sql . NullString { String : task . LiveConfirmedAt , Valid : task . LiveConfirmedAt != "" }
2026-08-09 22:55:14 +08:00
_ , err := q . Exec ( `
INSERT INTO tasks
2026-08-10 15:11:38 +08:00
(task_id, task_type, status, execution_mode, assigned_client,
2026-08-09 22:55:14 +08:00
syb_id, order_no, goods_id, shopee_sku_id,
pdd_goods_url, pdd_goods_id, pdd_options,
2026-08-10 23:10:28 +08:00
quantity, max_price_cent, live_confirmed_by, live_confirmed_at, created_by_user_id,
2026-08-10 15:11:38 +08:00
created_at, updated_at)
2026-08-10 23:10:28 +08:00
VALUES (?, 'purchase', 'assigned', ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)` ,
2026-08-10 15:11:38 +08:00
task . TaskID , task . ExecutionMode , task . AssignedClient , task . SybID , task . OrderNo ,
2026-08-09 22:55:14 +08:00
task . GoodsID , task . ShopeeSKUID , task . PddGoodsURL , task . PddGoodsID ,
2026-08-10 23:10:28 +08:00
task . PddOptions , task . Quantity , task . MaxPriceCent , liveConfirmedBy , liveConfirmedAt ,
sql . NullString { String : task . CreatedByUserID , Valid : strings . TrimSpace ( task . CreatedByUserID ) != "" }, now , now )
2026-08-09 22:55:14 +08:00
if err != nil {
return fmt . Errorf ( "创建顺运宝明细 %s 的采购任务失败: %w" , task . SybID , err )
}
return nil
}