diff --git a/server/app/goauto/access/purchaser.go b/server/app/goauto/access/purchaser.go index ac57c3a..acdeb8f 100644 --- a/server/app/goauto/access/purchaser.go +++ b/server/app/goauto/access/purchaser.go @@ -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}, diff --git a/server/app/goauto/migrations/migrate.go b/server/app/goauto/migrations/migrate.go index 29c9df8..86df8f6 100644 --- a/server/app/goauto/migrations/migrate.go +++ b/server/app/goauto/migrations/migrate.go @@ -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 { diff --git a/server/app/goauto/models/schema.go b/server/app/goauto/models/schema.go index 5a4e6e1..fe3a7ea 100644 --- a/server/app/goauto/models/schema.go +++ b/server/app/goauto/models/schema.go @@ -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: diff --git a/server/app/goauto/task/admin_handler.go b/server/app/goauto/task/admin_handler.go index 26fca9e..18acac3 100644 --- a/server/app/goauto/task/admin_handler.go +++ b/server/app/goauto/task/admin_handler.go @@ -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) diff --git a/server/app/goauto/task/agent_history.go b/server/app/goauto/task/agent_history.go index 7b851aa..454b2f0 100644 --- a/server/app/goauto/task/agent_history.go +++ b/server/app/goauto/task/agent_history.go @@ -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 diff --git a/server/app/goauto/task/cancel.go b/server/app/goauto/task/cancel.go new file mode 100644 index 0000000..e6a4253 --- /dev/null +++ b/server/app/goauto/task/cancel.go @@ -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 +} diff --git a/server/app/goauto/task/cancel_test.go b/server/app/goauto/task/cancel_test.go new file mode 100644 index 0000000..abb495d --- /dev/null +++ b/server/app/goauto/task/cancel_test.go @@ -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) + } +} diff --git a/server/app/goauto/task/router.go b/server/app/goauto/task/router.go index 4c828f8..f40e8a4 100644 --- a/server/app/goauto/task/router.go +++ b/server/app/goauto/task/router.go @@ -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) }