feat(server): 放开未开始采集任务的取消 (#297)
批量建单后无法中途叫停:删除只放行 failed,pending 一条都撤不掉,只能等设备 逐个跑完。图搜成功率并不高,批量越大越需要叫停。 新增 cancelled 终态,与 PurchaseTaskStatusCancelled 的既有约定对齐。 - pending → cancelled;running 拒绝取消。running 正在设备上操作 PDD,中途打断 后页面停在哪一步不可控,会影响下一个任务归位(#292 已为此付过代价)。语义是 「停止后续,当前这个跑完」。 - 取消与 Claim 的互斥点是同一条件更新。关键:Claim 领取时并不改 status,只写 device_id 和租约,因此条件里必须带 lease_expires_at,否则会把刚被领走的任务 误取消。 - 批量逐条更新、不包在一个事务里:一条因并发领取而跳过,不应回滚已成功取消的 其它任务。超出单批上限时以 HasMore 如实上报,不静默截断。 - status 的 CHECK 约束只认旧五值,GORM 在 MySQL 上不改写既有 CHECK,按同文件 ensureMySQLDirectSelectConstraint 的手法补幂等 DROP/ADD。cancelled 并入 syncGuardSlots 终态分支以满足 active_slot / device_run_slot 两个约束。 实施:sonnet 子代理,改动经独立复核与重跑验证。 Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01NTDbDcwbDw1TSAcE6wfh2F
This commit is contained in:
@@ -101,6 +101,8 @@ var AdminAPIs = []APIPermission{
|
||||
{"批量创建图搜采集任务", "/api/admin/v1/collection-tasks/image-search/batch", "POST", true},
|
||||
{"查看采集任务详情", "/api/admin/v1/collection-tasks/:taskId", "GET", true},
|
||||
{"重置采集任务", "/api/admin/v1/collection-tasks/:taskId/reset", "POST", true},
|
||||
{"取消采集任务", "/api/admin/v1/collection-tasks/:taskId/cancel", "POST", true},
|
||||
{"批量取消采集任务", "/api/admin/v1/collection-tasks/batch-cancel", "POST", true},
|
||||
{"删除采集任务", "/api/admin/v1/collection-tasks/:taskId", "DELETE", true},
|
||||
|
||||
{"查看采购任务", "/api/admin/v1/purchase-tasks", "GET", true},
|
||||
|
||||
@@ -101,7 +101,38 @@ func Migrate(db *gorm.DB) error {
|
||||
return err
|
||||
}
|
||||
}
|
||||
return ensureMySQLDirectSelectConstraint(db)
|
||||
if err := ensureMySQLDirectSelectConstraint(db); err != nil {
|
||||
return err
|
||||
}
|
||||
return ensureMySQLCollectionTaskCancelledConstraint(db)
|
||||
}
|
||||
|
||||
// ensureMySQLCollectionTaskCancelledConstraint 放开 collection_task.status 的
|
||||
// CHECK 约束以接受 'cancelled'(#297)。GORM 的 AutoMigrate 在 MySQL 上不会
|
||||
// 改写已存在的 CHECK 约束(同 ensureMySQLDirectSelectConstraint 的已知限制),
|
||||
// 所以新增取消状态必须像 direct_select 那次一样手动 DROP/ADD,否则老库的约束
|
||||
// 仍然只认旧的五个状态,取消写入会被数据库直接拒绝。
|
||||
func ensureMySQLCollectionTaskCancelledConstraint(db *gorm.DB) error {
|
||||
if db.Dialector.Name() != "mysql" {
|
||||
return nil
|
||||
}
|
||||
const name = "ck_collection_task_status"
|
||||
var constraints []struct {
|
||||
CheckClause string `gorm:"column:check_clause"`
|
||||
}
|
||||
if err := db.Raw(mysqlCheckConstraintQuery, "collection_task", name).Scan(&constraints).Error; err != nil {
|
||||
return err
|
||||
}
|
||||
if len(constraints) > 0 && strings.Contains(strings.ToLower(constraints[0].CheckClause), "cancelled") {
|
||||
return nil
|
||||
}
|
||||
if len(constraints) > 0 {
|
||||
if err := db.Exec("ALTER TABLE collection_task DROP CHECK " + name).Error; err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
return db.Exec("ALTER TABLE collection_task ADD CONSTRAINT " + name +
|
||||
" CHECK (status IN ('pending','running','completed','completed_partial','failed','cancelled'))").Error
|
||||
}
|
||||
|
||||
func ensureMySQLDirectSelectConstraint(db *gorm.DB) error {
|
||||
|
||||
@@ -17,6 +17,9 @@ const (
|
||||
TaskStatusCompleted = "completed"
|
||||
TaskStatusCompletedPartial = "completed_partial"
|
||||
TaskStatusFailed = "failed"
|
||||
// TaskStatusCancelled 与 PurchaseTaskStatusCancelled 的既有约定对齐(#297):
|
||||
// 只有 pending 任务允许取消,取消后是终态,不再被任何设备领取。
|
||||
TaskStatusCancelled = "cancelled"
|
||||
|
||||
CollectionTaskSourceAdmin = "admin"
|
||||
CollectionTaskSourceAgentCurrentPage = "agent_current_page"
|
||||
@@ -183,7 +186,7 @@ type CollectionTask struct {
|
||||
DeviceID *uint64 `json:"deviceId" gorm:"index;uniqueIndex:ux_collection_task_running_device,priority:1"`
|
||||
Device *AgentDevice `json:"-"`
|
||||
Source string `json:"source" gorm:"size:32;not null;default:admin;index;check:ck_collection_task_source,source IN ('admin','agent_current_page','image_search')"`
|
||||
Status string `json:"status" gorm:"size:24;not null;index;check:ck_collection_task_status,status IN ('pending','running','completed','completed_partial','failed')"`
|
||||
Status string `json:"status" gorm:"size:24;not null;index;check:ck_collection_task_status,status IN ('pending','running','completed','completed_partial','failed','cancelled')"`
|
||||
ActiveSlot *uint8 `json:"-" gorm:"uniqueIndex:ux_collection_task_active_product,priority:2;check:ck_collection_task_active_slot,(status IN ('pending','running') AND active_slot = 1) OR (status NOT IN ('pending','running') AND active_slot IS NULL)"`
|
||||
DeviceRunSlot *uint8 `json:"-" gorm:"uniqueIndex:ux_collection_task_running_device,priority:2;check:ck_collection_task_device_run_slot,(status = 'running' AND device_id IS NOT NULL AND device_run_slot = 1) OR (status <> 'running' AND device_run_slot IS NULL)"`
|
||||
URLSnapshot string `json:"urlSnapshot" gorm:"type:text;not null"`
|
||||
@@ -277,7 +280,7 @@ func (task *CollectionTask) syncGuardSlots() error {
|
||||
}
|
||||
task.ActiveSlot = &one
|
||||
task.DeviceRunSlot = &one
|
||||
case TaskStatusCompleted, TaskStatusCompletedPartial, TaskStatusFailed:
|
||||
case TaskStatusCompleted, TaskStatusCompletedPartial, TaskStatusFailed, TaskStatusCancelled:
|
||||
task.ActiveSlot = nil
|
||||
task.DeviceRunSlot = nil
|
||||
default:
|
||||
|
||||
@@ -161,6 +161,47 @@ func (handler Handler) AdminDelete(c *gin.Context) {
|
||||
c.JSON(http.StatusOK, gin.H{"code": 200, "data": response})
|
||||
}
|
||||
|
||||
func (handler Handler) AdminCancel(c *gin.Context) {
|
||||
id, err := taskID(c)
|
||||
if err != nil || id == 0 {
|
||||
writeError(c, serviceError("INVALID_REQUEST", "taskId 无效"))
|
||||
return
|
||||
}
|
||||
var request ActionRequest
|
||||
if err := decodeStrict(c, &request); err != nil {
|
||||
writeError(c, serviceError("INVALID_REQUEST", "请求 JSON 无效"))
|
||||
return
|
||||
}
|
||||
service, _, ok := handler.service(c)
|
||||
if !ok {
|
||||
return
|
||||
}
|
||||
response, err := service.Cancel(c.Request.Context(), id, request)
|
||||
if err != nil {
|
||||
writeError(c, err)
|
||||
return
|
||||
}
|
||||
c.JSON(http.StatusOK, gin.H{"code": 200, "data": response})
|
||||
}
|
||||
|
||||
func (handler Handler) AdminBatchCancel(c *gin.Context) {
|
||||
var request BatchCancelRequest
|
||||
if err := decodeStrict(c, &request); err != nil {
|
||||
writeError(c, serviceError("INVALID_REQUEST", "请求 JSON 无效"))
|
||||
return
|
||||
}
|
||||
service, _, ok := handler.service(c)
|
||||
if !ok {
|
||||
return
|
||||
}
|
||||
response, err := service.BatchCancel(c.Request.Context(), request)
|
||||
if err != nil {
|
||||
writeError(c, err)
|
||||
return
|
||||
}
|
||||
c.JSON(http.StatusOK, gin.H{"code": 200, "data": response})
|
||||
}
|
||||
|
||||
func decodeStrict(c *gin.Context, destination any) error {
|
||||
c.Request.Body = http.MaxBytesReader(c.Writer, c.Request.Body, 1<<20)
|
||||
decoder := json.NewDecoder(c.Request.Body)
|
||||
|
||||
@@ -238,7 +238,7 @@ func normalizeCollectionTaskNo(raw string) (uint64, bool, error) {
|
||||
|
||||
func validCollectionStatus(status string) bool {
|
||||
switch status {
|
||||
case models.TaskStatusPending, models.TaskStatusRunning, models.TaskStatusCompleted, models.TaskStatusCompletedPartial, models.TaskStatusFailed:
|
||||
case models.TaskStatusPending, models.TaskStatusRunning, models.TaskStatusCompleted, models.TaskStatusCompletedPartial, models.TaskStatusFailed, models.TaskStatusCancelled:
|
||||
return true
|
||||
default:
|
||||
return false
|
||||
|
||||
@@ -0,0 +1,163 @@
|
||||
package task
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"strings"
|
||||
|
||||
"go-admin/app/goauto/device"
|
||||
"go-admin/app/goauto/models"
|
||||
|
||||
"gorm.io/gorm"
|
||||
)
|
||||
|
||||
const maxBatchCancelItems = 500
|
||||
|
||||
// CancelResponse describes the outcome of cancelling one pending collection task.
|
||||
type CancelResponse struct {
|
||||
TaskID uint64 `json:"taskId"`
|
||||
Cancelled bool `json:"cancelled"`
|
||||
}
|
||||
|
||||
// BatchCancelRequest scopes a batch cancel by source and/or status. Only
|
||||
// pending tasks within the scope are actually cancelled; everything else in
|
||||
// scope is reported as skipped rather than causing the whole call to fail.
|
||||
type BatchCancelRequest struct {
|
||||
Source string `json:"source,omitempty"`
|
||||
Status string `json:"status,omitempty"`
|
||||
}
|
||||
|
||||
type BatchCancelSkippedItem struct {
|
||||
TaskID uint64 `json:"taskId"`
|
||||
Status string `json:"status"`
|
||||
}
|
||||
|
||||
type BatchCancelResponse struct {
|
||||
CancelledCount int `json:"cancelledCount"`
|
||||
CancelledIDs []uint64 `json:"cancelledIds"`
|
||||
Skipped []BatchCancelSkippedItem `json:"skipped"`
|
||||
// HasMore 表示范围内还有超出单批上限、本次未处理的任务。
|
||||
//
|
||||
// `[必须]` 必须如实上报,不能静默截断。调用方(采购员界面)会显示待取消
|
||||
// 的实时数量,截断而不告知会让「取消了 N 个」与看到的数字对不上,采购员
|
||||
// 以为全停了,实际上后面还在继续跑。
|
||||
HasMore bool `json:"hasMore"`
|
||||
}
|
||||
|
||||
// Cancel stops one pending collection task so a batch build-out can be
|
||||
// interrupted mid-flight (#297). Only pending tasks can be cancelled: a
|
||||
// running task is mid-operation on the phone and cutting it off leaves the
|
||||
// device's page position unknown, which is exactly the "手机停在深层页面"
|
||||
// failure #292 already paid for once. The semantics are "stop what has not
|
||||
// started, let what is running finish".
|
||||
func (service *Service) Cancel(ctx context.Context, taskID uint64, request ActionRequest) (CancelResponse, error) {
|
||||
if err := validateAction(taskID, request); err != nil {
|
||||
return CancelResponse{}, err
|
||||
}
|
||||
var response CancelResponse
|
||||
err := service.DB.WithContext(ctx).Transaction(func(tx *gorm.DB) error {
|
||||
var record models.CollectionTask
|
||||
if err := tx.First(&record, taskID).Error; err != nil {
|
||||
if errors.Is(err, gorm.ErrRecordNotFound) {
|
||||
return serviceError(CodeTaskNotFound, "任务不存在")
|
||||
}
|
||||
return internalError(err)
|
||||
}
|
||||
if record.Status == models.TaskStatusCancelled {
|
||||
response = CancelResponse{TaskID: taskID, Cancelled: true}
|
||||
return nil
|
||||
}
|
||||
if record.Status == models.TaskStatusRunning {
|
||||
return serviceError(CodeTaskStateConflict, "执行中的任务不能取消,请等待任务结束")
|
||||
}
|
||||
if record.Status != models.TaskStatusPending {
|
||||
return serviceError(CodeTaskStateConflict, "只有待执行任务可以取消")
|
||||
}
|
||||
// `[必须]` 取消与设备领取(Claim)必须互斥,且不能靠"先读上面的
|
||||
// record.Status 再写"来判断——那一读一写之间设备完全可能已经把任务
|
||||
// 领走。真正的互斥点是这条条件更新,和 service.go 里 Claim 用的是
|
||||
// 同一模式:WHERE status = pending AND lease 未生效,并校验
|
||||
// RowsAffected。这里必须带 lease_expires_at 条件,因为 Claim 领取
|
||||
// 任务时并不改 status(仍是 pending,只是设了 device_id 和租约),
|
||||
// 只看 status 会把"已被领取、马上要开始"的任务误判成可取消。
|
||||
// RowsAffected 为 0 说明任务已被领取或状态已变化,取消必须失败。
|
||||
now := service.Now()
|
||||
result := tx.Session(&gorm.Session{SkipHooks: true}).Model(&models.CollectionTask{}).
|
||||
Where("id = ? AND status = ? AND (lease_expires_at IS NULL OR lease_expires_at <= ?)", record.ID, models.TaskStatusPending, now).
|
||||
Updates(map[string]any{
|
||||
"status": models.TaskStatusCancelled, "active_slot": nil, "device_run_slot": nil,
|
||||
})
|
||||
if result.Error != nil {
|
||||
return internalError(result.Error)
|
||||
}
|
||||
if result.RowsAffected != 1 {
|
||||
return serviceError(CodeTaskStateConflict, "任务状态已变化,可能已被设备领取")
|
||||
}
|
||||
response = CancelResponse{TaskID: taskID, Cancelled: true}
|
||||
return nil
|
||||
})
|
||||
return response, err
|
||||
}
|
||||
|
||||
// BatchCancel cancels every pending task within the requested source/status
|
||||
// scope. It never touches running, completed, failed or already-cancelled
|
||||
// tasks, and each row is cancelled through the same conditional update Cancel
|
||||
// uses, so a task claimed by a device between the scope query and the update
|
||||
// is safely skipped instead of being cancelled out from under it.
|
||||
//
|
||||
// `[必须]` 逐条更新、不包在一个事务里是刻意的:一条因并发领取而跳过,不应该
|
||||
// 回滚已经成功取消的其它任务。采购员要的是「能停多少停多少」,而不是全有或
|
||||
// 全无——批量越大,中途有一条被领走的概率越高。
|
||||
func (service *Service) BatchCancel(ctx context.Context, request BatchCancelRequest) (BatchCancelResponse, error) {
|
||||
source := strings.TrimSpace(request.Source)
|
||||
if source != "" && source != models.CollectionTaskSourceAdmin && source != models.CollectionTaskSourceAgentCurrentPage && source != models.CollectionTaskSourceImageSearch {
|
||||
return BatchCancelResponse{}, serviceError(device.CodeInvalidRequest, "source 无效")
|
||||
}
|
||||
status := strings.TrimSpace(request.Status)
|
||||
if status != "" && !validCollectionStatus(status) {
|
||||
return BatchCancelResponse{}, serviceError(device.CodeInvalidRequest, "status 无效")
|
||||
}
|
||||
response := BatchCancelResponse{CancelledIDs: []uint64{}, Skipped: []BatchCancelSkippedItem{}}
|
||||
query := service.DB.WithContext(ctx).Model(&models.CollectionTask{})
|
||||
if source != "" {
|
||||
query = query.Where("source = ?", source)
|
||||
}
|
||||
if status != "" {
|
||||
query = query.Where("status = ?", status)
|
||||
}
|
||||
var candidates []models.CollectionTask
|
||||
// 多取一条用于判断范围内是否还有未处理的任务,多出来的那条不参与取消。
|
||||
if err := query.Order("id ASC").Limit(maxBatchCancelItems + 1).Find(&candidates).Error; err != nil {
|
||||
return BatchCancelResponse{}, internalError(err)
|
||||
}
|
||||
if len(candidates) > maxBatchCancelItems {
|
||||
candidates = candidates[:maxBatchCancelItems]
|
||||
response.HasMore = true
|
||||
}
|
||||
now := service.Now()
|
||||
for _, candidate := range candidates {
|
||||
if candidate.Status != models.TaskStatusPending {
|
||||
response.Skipped = append(response.Skipped, BatchCancelSkippedItem{TaskID: candidate.ID, Status: candidate.Status})
|
||||
continue
|
||||
}
|
||||
// `[必须]` 与 Cancel 相同的互斥条件:不能只看 status,还要排除已被
|
||||
// Claim 生效租约(未过期)的任务,否则批量取消会把刚被设备领走、
|
||||
// 状态仍是 pending 的任务连带取消掉。
|
||||
result := service.DB.WithContext(ctx).Session(&gorm.Session{SkipHooks: true}).Model(&models.CollectionTask{}).
|
||||
Where("id = ? AND status = ? AND (lease_expires_at IS NULL OR lease_expires_at <= ?)", candidate.ID, models.TaskStatusPending, now).
|
||||
Updates(map[string]any{
|
||||
"status": models.TaskStatusCancelled, "active_slot": nil, "device_run_slot": nil,
|
||||
})
|
||||
if result.Error != nil {
|
||||
return BatchCancelResponse{}, internalError(result.Error)
|
||||
}
|
||||
if result.RowsAffected != 1 {
|
||||
// 领取发生在查询候选和这次条件更新之间;不是误取消,而是如实上报跳过。
|
||||
response.Skipped = append(response.Skipped, BatchCancelSkippedItem{TaskID: candidate.ID, Status: models.TaskStatusPending})
|
||||
continue
|
||||
}
|
||||
response.CancelledCount++
|
||||
response.CancelledIDs = append(response.CancelledIDs, candidate.ID)
|
||||
}
|
||||
return response, nil
|
||||
}
|
||||
@@ -0,0 +1,263 @@
|
||||
package task
|
||||
|
||||
import (
|
||||
"context"
|
||||
"testing"
|
||||
|
||||
"go-admin/app/goauto/models"
|
||||
|
||||
"github.com/google/uuid"
|
||||
)
|
||||
|
||||
// TestCancelPendingTaskBlocksFutureClaim 验证 #297 的核心验收点:pending 任务
|
||||
// 可以取消,取消后这个任务不再能被任何设备领取。
|
||||
func TestCancelPendingTaskBlocksFutureClaim(t *testing.T) {
|
||||
db := openTaskDatabase(t)
|
||||
service := newTaskService(db)
|
||||
_, token := registerTaskDevice(t, db, "device-one")
|
||||
record := createTask(t, db, nil)
|
||||
|
||||
response, err := service.Cancel(context.Background(), record.ID, ActionRequest{RequestID: uuid.NewString()})
|
||||
if err != nil {
|
||||
t.Fatalf("cancel: %v", err)
|
||||
}
|
||||
if !response.Cancelled {
|
||||
t.Fatalf("expected cancelled=true, got %+v", response)
|
||||
}
|
||||
|
||||
var stored models.CollectionTask
|
||||
if err := db.First(&stored, record.ID).Error; err != nil {
|
||||
t.Fatalf("load task: %v", err)
|
||||
}
|
||||
if stored.Status != models.TaskStatusCancelled {
|
||||
t.Fatalf("expected status cancelled, got %q", stored.Status)
|
||||
}
|
||||
if stored.ActiveSlot != nil || stored.DeviceRunSlot != nil {
|
||||
t.Fatalf("expected guard slots cleared on cancelled task, got active=%v run=%v", stored.ActiveSlot, stored.DeviceRunSlot)
|
||||
}
|
||||
|
||||
if _, err := service.Claim(context.Background(), record.ID, ActionRequest{RequestID: uuid.NewString()}, token); err == nil {
|
||||
t.Fatalf("expected claim of cancelled task to fail")
|
||||
} else if code := taskErrorCode(t, err); code != CodeTaskStateConflict {
|
||||
t.Fatalf("expected TASK_STATE_CONFLICT, got %s", code)
|
||||
}
|
||||
}
|
||||
|
||||
// TestCancelRunningTaskRejected 验证 running 任务禁止取消,错误可读(#292 教训:
|
||||
// 设备正在操作拼多多,中途打断会让手机停在不可控页面)。
|
||||
func TestCancelRunningTaskRejected(t *testing.T) {
|
||||
db := openTaskDatabase(t)
|
||||
service := newTaskService(db)
|
||||
deviceRecord, _ := registerTaskDevice(t, db, "device-one")
|
||||
record := createTask(t, db, &deviceRecord.ID)
|
||||
if err := record.SetStatus(models.TaskStatusRunning); err != nil {
|
||||
t.Fatalf("set status: %v", err)
|
||||
}
|
||||
if err := db.Model(&record).Updates(map[string]any{"status": record.Status, "active_slot": record.ActiveSlot, "device_run_slot": record.DeviceRunSlot}).Error; err != nil {
|
||||
t.Fatalf("persist running status: %v", err)
|
||||
}
|
||||
|
||||
_, err := service.Cancel(context.Background(), record.ID, ActionRequest{RequestID: uuid.NewString()})
|
||||
if err == nil {
|
||||
t.Fatalf("expected cancel of running task to fail")
|
||||
}
|
||||
if code := taskErrorCode(t, err); code != CodeTaskStateConflict {
|
||||
t.Fatalf("expected TASK_STATE_CONFLICT, got %s", code)
|
||||
}
|
||||
|
||||
var stored models.CollectionTask
|
||||
if err := db.First(&stored, record.ID).Error; err != nil {
|
||||
t.Fatalf("load task: %v", err)
|
||||
}
|
||||
if stored.Status != models.TaskStatusRunning {
|
||||
t.Fatalf("running task status must be untouched, got %q", stored.Status)
|
||||
}
|
||||
}
|
||||
|
||||
// TestCancelAfterClaimFailsInsteadOfMiscancelling 是并发验收点:任务已被设备
|
||||
// 领取(不再是 pending)之后再取消,必须失败而不是误取消。取消用的条件更新
|
||||
// (WHERE status = pending + RowsAffected)和 Claim 是同一互斥模式,这里直接
|
||||
// 复现"先领取,后取消"的时序来证明它真的互斥。
|
||||
func TestCancelAfterClaimFailsInsteadOfMiscancelling(t *testing.T) {
|
||||
db := openTaskDatabase(t)
|
||||
service := newTaskService(db)
|
||||
deviceRecord, token := registerTaskDevice(t, db, "device-one")
|
||||
record := createTask(t, db, &deviceRecord.ID)
|
||||
|
||||
if _, err := service.Claim(context.Background(), record.ID, ActionRequest{RequestID: uuid.NewString()}, token); err != nil {
|
||||
t.Fatalf("claim: %v", err)
|
||||
}
|
||||
|
||||
_, err := service.Cancel(context.Background(), record.ID, ActionRequest{RequestID: uuid.NewString()})
|
||||
if err == nil {
|
||||
t.Fatalf("expected cancel to fail once the task has been claimed")
|
||||
}
|
||||
if code := taskErrorCode(t, err); code != CodeTaskStateConflict {
|
||||
t.Fatalf("expected TASK_STATE_CONFLICT, got %s", code)
|
||||
}
|
||||
|
||||
var stored models.CollectionTask
|
||||
if err := db.First(&stored, record.ID).Error; err != nil {
|
||||
t.Fatalf("load task: %v", err)
|
||||
}
|
||||
if stored.Status != models.TaskStatusPending {
|
||||
t.Fatalf("claimed task must remain pending (claim keeps status pending until start), got %q", stored.Status)
|
||||
}
|
||||
if stored.DeviceID == nil || *stored.DeviceID != deviceRecord.ID {
|
||||
t.Fatalf("claimed task must keep its device assignment, got %+v", stored.DeviceID)
|
||||
}
|
||||
}
|
||||
|
||||
func TestCancelNotFoundAndAlreadyCancelledIsIdempotent(t *testing.T) {
|
||||
db := openTaskDatabase(t)
|
||||
service := newTaskService(db)
|
||||
|
||||
if _, err := service.Cancel(context.Background(), 999999, ActionRequest{RequestID: uuid.NewString()}); err == nil {
|
||||
t.Fatalf("expected not found error")
|
||||
} else if code := taskErrorCode(t, err); code != CodeTaskNotFound {
|
||||
t.Fatalf("expected TASK_NOT_FOUND, got %s", code)
|
||||
}
|
||||
|
||||
record := createTask(t, db, nil)
|
||||
if _, err := service.Cancel(context.Background(), record.ID, ActionRequest{RequestID: uuid.NewString()}); err != nil {
|
||||
t.Fatalf("first cancel: %v", err)
|
||||
}
|
||||
response, err := service.Cancel(context.Background(), record.ID, ActionRequest{RequestID: uuid.NewString()})
|
||||
if err != nil {
|
||||
t.Fatalf("second cancel on already-cancelled task should be idempotent, got error: %v", err)
|
||||
}
|
||||
if !response.Cancelled {
|
||||
t.Fatalf("expected cancelled=true on replay, got %+v", response)
|
||||
}
|
||||
}
|
||||
|
||||
func TestBatchCancelScopesBySourceAndStatus(t *testing.T) {
|
||||
db := openTaskDatabase(t)
|
||||
service := newTaskService(db)
|
||||
|
||||
pendingAdmin := createTask(t, db, nil)
|
||||
otherPendingAdmin := createTask(t, db, nil)
|
||||
|
||||
imageSearchTask := createTask(t, db, nil)
|
||||
if err := db.Model(&imageSearchTask).Update("source", models.CollectionTaskSourceImageSearch).Error; err != nil {
|
||||
t.Fatalf("set source: %v", err)
|
||||
}
|
||||
|
||||
completedTask := createTask(t, db, nil)
|
||||
if err := db.Model(&completedTask).Updates(map[string]any{"status": models.TaskStatusCompleted, "active_slot": nil, "device_run_slot": nil}).Error; err != nil {
|
||||
t.Fatalf("set completed: %v", err)
|
||||
}
|
||||
|
||||
response, err := service.BatchCancel(context.Background(), BatchCancelRequest{Source: models.CollectionTaskSourceAdmin, Status: models.TaskStatusPending})
|
||||
if err != nil {
|
||||
t.Fatalf("batch cancel: %v", err)
|
||||
}
|
||||
if response.CancelledCount != 2 {
|
||||
t.Fatalf("expected 2 cancelled, got %d (%+v)", response.CancelledCount, response)
|
||||
}
|
||||
cancelledSet := map[uint64]bool{}
|
||||
for _, id := range response.CancelledIDs {
|
||||
cancelledSet[id] = true
|
||||
}
|
||||
if !cancelledSet[pendingAdmin.ID] || !cancelledSet[otherPendingAdmin.ID] {
|
||||
t.Fatalf("expected both pending admin tasks cancelled, got %+v", response.CancelledIDs)
|
||||
}
|
||||
|
||||
// 范围外的任务必须原样不动:图搜任务未被 source=admin 命中,已完成任务不是 pending。
|
||||
var storedImageSearch, storedCompleted models.CollectionTask
|
||||
if err := db.First(&storedImageSearch, imageSearchTask.ID).Error; err != nil {
|
||||
t.Fatalf("load image search task: %v", err)
|
||||
}
|
||||
if storedImageSearch.Status != models.TaskStatusPending {
|
||||
t.Fatalf("image search task outside source scope must stay pending, got %q", storedImageSearch.Status)
|
||||
}
|
||||
if err := db.First(&storedCompleted, completedTask.ID).Error; err != nil {
|
||||
t.Fatalf("load completed task: %v", err)
|
||||
}
|
||||
if storedCompleted.Status != models.TaskStatusCompleted {
|
||||
t.Fatalf("completed task must be untouched, got %q", storedCompleted.Status)
|
||||
}
|
||||
}
|
||||
|
||||
// TestBatchCancelDoesNotAffectOtherTasksOfSameProduct 验证取消不影响同商品
|
||||
// 其它任务:取消一个商品下的 pending 任务,不应波及该商品其它独立任务。
|
||||
func TestCancelDoesNotAffectOtherTaskOfSameProduct(t *testing.T) {
|
||||
db := openTaskDatabase(t)
|
||||
service := newTaskService(db)
|
||||
|
||||
goodsID := uuid.NewString()
|
||||
product := models.PDDProduct{GoodsID: goodsID, URL: "https://mobile.yangkeduo.com/goods.html?goods_id=" + goodsID}
|
||||
if err := db.Create(&product).Error; err != nil {
|
||||
t.Fatalf("create product: %v", err)
|
||||
}
|
||||
rule := models.CollectionRule{Name: "rule", ContentJSON: `{"steps":[]}`}
|
||||
if err := db.Create(&rule).Error; err != nil {
|
||||
t.Fatalf("create rule: %v", err)
|
||||
}
|
||||
pending := models.CollectionTask{PDDProductID: &product.ID, RuleID: rule.ID, Status: models.TaskStatusPending, URLSnapshot: product.URL, GoodsIDSnapshot: product.GoodsID, RuleSnapshot: rule.ContentJSON}
|
||||
if err := db.Create(&pending).Error; err != nil {
|
||||
t.Fatalf("create pending task: %v", err)
|
||||
}
|
||||
failed := models.CollectionTask{PDDProductID: &product.ID, RuleID: rule.ID, Status: models.TaskStatusFailed, URLSnapshot: product.URL, GoodsIDSnapshot: product.GoodsID, RuleSnapshot: rule.ContentJSON}
|
||||
if err := db.Create(&failed).Error; err != nil {
|
||||
t.Fatalf("create failed task: %v", err)
|
||||
}
|
||||
|
||||
if _, err := service.Cancel(context.Background(), pending.ID, ActionRequest{RequestID: uuid.NewString()}); err != nil {
|
||||
t.Fatalf("cancel: %v", err)
|
||||
}
|
||||
|
||||
var storedFailed models.CollectionTask
|
||||
if err := db.First(&storedFailed, failed.ID).Error; err != nil {
|
||||
t.Fatalf("load failed task: %v", err)
|
||||
}
|
||||
if storedFailed.Status != models.TaskStatusFailed {
|
||||
t.Fatalf("failed task of same product must be untouched, got %q", storedFailed.Status)
|
||||
}
|
||||
}
|
||||
|
||||
// `[必须]` 超出单批上限时必须如实上报,不能静默截断。采购员界面会显示待取消的
|
||||
// 实时数量,截断而不告知会让「取消了 N 个」和眼前的数字对不上,人以为全停了,
|
||||
// 实际后面还在继续跑。
|
||||
func TestBatchCancelReportsWhenScopeExceedsTheLimit(t *testing.T) {
|
||||
db := openTaskDatabase(t)
|
||||
service := newTaskService(db)
|
||||
|
||||
for i := 0; i < maxBatchCancelItems+3; i++ {
|
||||
createTask(t, db, nil)
|
||||
}
|
||||
|
||||
response, err := service.BatchCancel(context.Background(), BatchCancelRequest{Status: models.TaskStatusPending})
|
||||
if err != nil {
|
||||
t.Fatalf("batch cancel: %v", err)
|
||||
}
|
||||
if response.CancelledCount != maxBatchCancelItems {
|
||||
t.Fatalf("cancelled %d, want the batch limit %d", response.CancelledCount, maxBatchCancelItems)
|
||||
}
|
||||
if !response.HasMore {
|
||||
t.Fatal("范围内还有未处理的任务,HasMore 必须为 true")
|
||||
}
|
||||
|
||||
// 再来一次应当把剩下的收掉,并且不再报 HasMore。
|
||||
rest, err := service.BatchCancel(context.Background(), BatchCancelRequest{Status: models.TaskStatusPending})
|
||||
if err != nil {
|
||||
t.Fatalf("batch cancel rest: %v", err)
|
||||
}
|
||||
if rest.CancelledCount != 3 || rest.HasMore {
|
||||
t.Fatalf("second pass cancelled=%d hasMore=%v, want 3 / false", rest.CancelledCount, rest.HasMore)
|
||||
}
|
||||
}
|
||||
|
||||
func TestBatchCancelWithinLimitDoesNotReportMore(t *testing.T) {
|
||||
db := openTaskDatabase(t)
|
||||
service := newTaskService(db)
|
||||
createTask(t, db, nil)
|
||||
|
||||
response, err := service.BatchCancel(context.Background(), BatchCancelRequest{Status: models.TaskStatusPending})
|
||||
if err != nil {
|
||||
t.Fatalf("batch cancel: %v", err)
|
||||
}
|
||||
if response.CancelledCount != 1 || response.HasMore {
|
||||
t.Fatalf("cancelled=%d hasMore=%v, want 1 / false", response.CancelledCount, response.HasMore)
|
||||
}
|
||||
}
|
||||
@@ -32,7 +32,9 @@ func InitRouter(engine *gin.Engine, auth *jwt.GinJWTMiddleware) {
|
||||
admin.POST("", handler.AdminCreate)
|
||||
admin.POST("/batch", handler.AdminBatchCreate)
|
||||
admin.POST("/image-search/batch", handler.AdminBatchCreateImageSearch)
|
||||
admin.POST("/batch-cancel", handler.AdminBatchCancel)
|
||||
admin.GET("/:taskId", handler.AdminDetail)
|
||||
admin.POST("/:taskId/reset", handler.AdminReset)
|
||||
admin.POST("/:taskId/cancel", handler.AdminCancel)
|
||||
admin.DELETE("/:taskId", handler.AdminDelete)
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user