From 7b6cc00503d3a685a3271d0512e9ddaaa3a4a597 Mon Sep 17 00:00:00 2001 From: QiuSW <105186638@qq.com> Date: Wed, 26 Aug 2026 10:17:45 +0800 Subject: [PATCH] feat(agent): implement #91 controlled recollection --- android/app/build.gradle.kts | 4 +- .../goauto/agent/network/AgentApiClient.kt | 8 ++ .../goauto/agent/ui/TaskHistoryFragment.kt | 65 +++++++++ .../goauto/agent/CollectionResetPolicyTest.kt | 23 ++++ docs/02-architecture-and-code-map.md | 10 +- docs/03-business-rules-and-glossary.md | 12 +- docs/08-agent-api-contract.md | 21 ++- server/app/goauto/task/agent_reset_test.go | 124 ++++++++++++++++++ server/app/goauto/task/handler.go | 24 ++++ server/app/goauto/task/lifecycle_service.go | 65 +++++++++ server/app/goauto/task/router.go | 1 + 11 files changed, 349 insertions(+), 8 deletions(-) create mode 100644 android/app/src/test/java/cn/ilapage/goauto/agent/CollectionResetPolicyTest.kt create mode 100644 server/app/goauto/task/agent_reset_test.go diff --git a/android/app/build.gradle.kts b/android/app/build.gradle.kts index 347db34..eeab937 100644 --- a/android/app/build.gradle.kts +++ b/android/app/build.gradle.kts @@ -11,8 +11,8 @@ android { applicationId = "cn.ilapage.goauto.agent" minSdk = 23 targetSdk = 34 - versionCode = 6 - versionName = "0.4.1" + versionCode = 7 + versionName = "0.5.0" testInstrumentationRunner = "androidx.test.runner.AndroidJUnitRunner" diff --git a/android/app/src/main/java/cn/ilapage/goauto/agent/network/AgentApiClient.kt b/android/app/src/main/java/cn/ilapage/goauto/agent/network/AgentApiClient.kt index d1f3bcb..1f1c87f 100644 --- a/android/app/src/main/java/cn/ilapage/goauto/agent/network/AgentApiClient.kt +++ b/android/app/src/main/java/cn/ilapage/goauto/agent/network/AgentApiClient.kt @@ -92,6 +92,8 @@ data class CollectionHistoryDetail( val missing: List, ) +data class CollectionResetResult(val taskId: Long, val status: String, val replayed: Boolean) + data class PurchaseHistoryItem( val taskId: Long, val status: String, @@ -242,6 +244,12 @@ class AgentApiClient(private val serverUrl: String) { ) } + fun resetCollectionTask(taskId: Long, requestId: String, token: String): CollectionResetResult { + val payload = JSONObject().put("requestId", requestId) + val data = requireNotNull(request("POST", "/api/agent/v1/collection-tasks/$taskId/reset", payload, token)).getJSONObject("data") + return CollectionResetResult(data.getLong("taskId"), data.getString("status"), data.optBoolean("replayed")) + } + fun purchaseHistory(token: String, page: Int, status: String?, taskNo: String?): HistoryPage { val data = requireNotNull(request("GET", historyPath("/api/agent/v1/purchase-tasks", page, status, taskNo), null, token)).getJSONObject("data") return HistoryPage( diff --git a/android/app/src/main/java/cn/ilapage/goauto/agent/ui/TaskHistoryFragment.kt b/android/app/src/main/java/cn/ilapage/goauto/agent/ui/TaskHistoryFragment.kt index b4fde74..ff07718 100644 --- a/android/app/src/main/java/cn/ilapage/goauto/agent/ui/TaskHistoryFragment.kt +++ b/android/app/src/main/java/cn/ilapage/goauto/agent/ui/TaskHistoryFragment.kt @@ -20,9 +20,12 @@ import cn.ilapage.goauto.agent.network.CollectionHistoryItem import cn.ilapage.goauto.agent.network.HistoryPage import cn.ilapage.goauto.agent.network.PurchaseHistoryDetail import cn.ilapage.goauto.agent.network.PurchaseHistoryItem +import cn.ilapage.goauto.agent.service.AgentForegroundService import cn.ilapage.goauto.agent.service.AgentSettingsStore import com.google.android.material.button.MaterialButton +import com.google.android.material.dialog.MaterialAlertDialogBuilder import java.util.Locale +import java.util.UUID internal object TaskNumberParser { fun normalize(raw: String, collection: Boolean): String? { @@ -40,6 +43,18 @@ internal object TaskNumberParser { } } +internal object CollectionResetPolicy { + private val terminalStatuses = setOf("completed", "completed_partial", "failed") + + fun isAllowed(status: String): Boolean = status in terminalStatuses + + fun confirmationMessage(status: String): String = if (status == "failed") { + "重新采集会清除本任务现有的错误和采集结果,然后重新加入当前设备的任务队列。" + } else { + "重新采集会清除本任务现有的全部采集结果,然后重新加入当前设备的任务队列。" + } +} + class TaskHistoryFragment : Fragment() { private val collection: Boolean get() = requireArguments().getBoolean(ARG_COLLECTION) private lateinit var pageColumn: LinearLayout @@ -291,6 +306,56 @@ class TaskHistoryFragment : Fragment() { }), collectionCardParams()) if (detail.missing.isNotEmpty()) resultColumn.addView(context.centeredMessage("缺失项", detail.missing.joinToString("、"))) if (task.errorMessage != null) resultColumn.addView(context.centeredMessage(task.errorMessage, "错误代码:${task.errorCode ?: "—"}")) + if (CollectionResetPolicy.isAllowed(task.status)) { + resultColumn.addView(MaterialButton(context).apply { + text = "重新采集" + minimumHeight = context.dp(48) + contentDescription = "重新采集任务 ${task.taskId}" + setOnClickListener { confirmReset(task) } + }, collectionCardParams()) + } + } + + private fun confirmReset(task: CollectionHistoryItem) { + MaterialAlertDialogBuilder(requireContext()) + .setTitle("重新采集 #${task.taskId}?") + .setMessage(CollectionResetPolicy.confirmationMessage(task.status)) + .setNegativeButton("取消", null) + .setPositiveButton("确认重新采集") { _, _ -> resetCollectionTask(task.taskId) } + .show() + } + + private fun resetCollectionTask(taskId: Long) { + val context = requireContext() + val credentials = runCatching { SecureDeviceStore(context).credentials() }.getOrNull() + val serverUrl = AgentSettingsStore(context).serverUrl() + if (credentials == null || serverUrl.isBlank()) { + showMessage("无法重新采集", "设备尚未连接服务端,请先检查设置。", "返回任务详情") { loadCollectionDetail(taskId) } + return + } + val generation = ++requestGeneration + showLoading("正在加入重新采集队列…") + Thread { + runCatching { AgentApiClient(serverUrl).resetCollectionTask(taskId, UUID.randomUUID().toString(), credentials.token) } + .onSuccess { result -> + resultColumn.post { + if (!isAdded || generation != requestGeneration) return@post + AgentForegroundService.start(requireContext()) + showMessage( + "已加入重新采集队列", + "任务 #${result.taskId} 将由当前设备通过正常调度重新执行。", + "返回采集记录", + ) { page = 1; load() } + } + } + .onFailure { error -> + resultColumn.post { + if (isAdded && generation == requestGeneration) { + showMessage("无法重新采集", error.message ?: "请检查设备状态后重试。", "返回任务详情") { loadCollectionDetail(taskId) } + } + } + } + }.start() } private fun renderPurchaseDetail(detail: PurchaseHistoryDetail) { diff --git a/android/app/src/test/java/cn/ilapage/goauto/agent/CollectionResetPolicyTest.kt b/android/app/src/test/java/cn/ilapage/goauto/agent/CollectionResetPolicyTest.kt new file mode 100644 index 0000000..6a2b969 --- /dev/null +++ b/android/app/src/test/java/cn/ilapage/goauto/agent/CollectionResetPolicyTest.kt @@ -0,0 +1,23 @@ +package cn.ilapage.goauto.agent + +import cn.ilapage.goauto.agent.ui.CollectionResetPolicy +import org.junit.Assert.assertFalse +import org.junit.Assert.assertTrue +import org.junit.Test + +class CollectionResetPolicyTest { + @Test + fun `only collection terminal states can reset`() { + assertTrue(CollectionResetPolicy.isAllowed("completed")) + assertTrue(CollectionResetPolicy.isAllowed("completed_partial")) + assertTrue(CollectionResetPolicy.isAllowed("failed")) + assertFalse(CollectionResetPolicy.isAllowed("pending")) + assertFalse(CollectionResetPolicy.isAllowed("running")) + } + + @Test + fun `confirmation always explains result deletion`() { + assertTrue(CollectionResetPolicy.confirmationMessage("failed").contains("清除")) + assertTrue(CollectionResetPolicy.confirmationMessage("completed").contains("全部采集结果")) + } +} diff --git a/docs/02-architecture-and-code-map.md b/docs/02-architecture-and-code-map.md index 9d830c0..1ce8e46 100644 --- a/docs/02-architecture-and-code-map.md +++ b/docs/02-architecture-and-code-map.md @@ -2,8 +2,8 @@ generated: true (请先修改 Gitea Wiki,禁止直接编辑本文件) wiki_page: Architecture-and-Code-Map wiki_url: https://git.ilapage.cn/OPC/goauto/wiki/Architecture-and-Code-Map.- -wiki_revision: 3b7a0cad581f615124aeeee8f8d3baf83e1dca28 -synchronized_at: 2026-08-26T01:17:16Z +wiki_revision: 9a52b45dc154e741ee2abcd9e74b947fa3745a28 +synchronized_at: 2026-08-26T02:11:15Z # 架构与代码地图 @@ -173,3 +173,9 @@ Android Portal/Agent - `AgentApiClient` 调用服务端当前设备历史接口;任务编号规范化由客户端先处理,服务端再次校验。 - 服务端采集与采购模块各自使用 `agent_history.go` 负责 Device Token 设备隔离、最近 30 天过滤、分页和最小化 DTO。 - `purchase_task.actual_unit_price_cent` 是可空字段,只记录 Agent 实际观察到的采购单价;旧任务不回填。 + +## Agent 受控重新采集(#91) + +- 服务端 `task/lifecycle_service.go` 的同一 Reset 事务同时服务管理端和设备端;设备端入口使用显式 Device Token 模式,不能用空 Token 退化为管理端路径。 +- `POST /api/agent/v1/collection-tasks/{taskId}/reset` 只返回最小状态,设备归属、忙碌和商品冲突在事务内检查。 +- Android `TaskHistoryFragment` 只在采集终态详情显示“重新采集”,确认成功后调用 `AgentForegroundService.start` 唤醒现有轮询,任务执行仍由 `next/claim/start` 完成。 diff --git a/docs/03-business-rules-and-glossary.md b/docs/03-business-rules-and-glossary.md index 99e10f8..5b8f85d 100644 --- a/docs/03-business-rules-and-glossary.md +++ b/docs/03-business-rules-and-glossary.md @@ -2,8 +2,8 @@ generated: true (请先修改 Gitea Wiki,禁止直接编辑本文件) wiki_page: Business-Rules-and-Glossary wiki_url: https://git.ilapage.cn/OPC/goauto/wiki/Business-Rules-and-Glossary.- -wiki_revision: c956def096806541d5060c771504928d121eb127 -synchronized_at: 2026-08-26T01:17:27Z +wiki_revision: e83a27ca331f23a9219fbf4bde9020104192762a +synchronized_at: 2026-08-26T02:11:24Z # 业务规则与术语 @@ -231,3 +231,11 @@ synchronized_at: 2026-08-26T01:17:27Z - 采集历史可以查看结构化采集结果和错误;采购历史可以查看结构化采购数据和错误,但全部为只读。 - 历史接口遵循数据最小化:不下发 Device Token、规则快照、PDD URL、收货地址、原始控件树或截图。 - 实际采购单价只在 Agent 确实观察到时保存;旧任务或未观察到价格时显示“未记录”,不推算、不回填。 + +## Agent 受控重新采集 + +- 当前设备可以从自身采集任务详情重新采集已完成、部分完成或失败任务;采购任务始终只读,不能在 Agent 重新采购。 +- 重新采集属于破坏性状态操作,必须先显示确认框并明确提示会清除当前结构化结果和错误;用户取消时不得请求服务端。 +- 服务端是最终事实来源:校验 Device Token、设备归属、在线/空闲、任务终态和同商品活动任务冲突,Android 不得本地绕过。 +- 重置成功后复用原任务,保留 URL、goods_id、设备与规则快照,清除旧结果后恢复 `pending`;执行仍走正常租约和串行调度。 +- 不支持离线排队、跨设备重置、自动换机、重新采购、创建订单或支付。 diff --git a/docs/08-agent-api-contract.md b/docs/08-agent-api-contract.md index ec7d6ec..95042c4 100644 --- a/docs/08-agent-api-contract.md +++ b/docs/08-agent-api-contract.md @@ -2,8 +2,8 @@ generated: true (请先修改 Gitea Wiki,禁止直接编辑本文件) wiki_page: Android-Agent-API-Contract wiki_url: https://git.ilapage.cn/OPC/goauto/wiki/Android-Agent-API-Contract.- -wiki_revision: 0f83e8f7333f2bec36f7b42f9ab558ec3b1e265e -synchronized_at: 2026-08-26T01:18:14Z +wiki_revision: 7f48742d36b23e7d2c01d83f714812b810bfc34f +synchronized_at: 2026-08-26T02:12:04Z # MVP 共享 API 契约 @@ -550,3 +550,20 @@ GET /api/agent/v1/purchase-tasks/{taskId} - 采集详情返回标题、店铺、销量、评价数、规格维度、颜色价格、SKU、缺失项和结构化错误。 - 采购详情返回蝦皮订单号、PDD 商品、已匹配颜色/尺码、数量、实际单价、PDD 订单号、下单时间和结构化错误。 - Agent 提交采购结果时可携带 `actualUnitPriceCent`(人民币分,非负)。服务端只保存 Agent 实际观察到的值;历史任务或未观察到价格的结果保持 `null`,客户端显示“未记录”。 + +## Agent 受控重新采集(#91) + +```http +POST /api/agent/v1/collection-tasks/{taskId}/reset +Authorization: Bearer +Content-Type: application/json + +{"requestId":""} +``` + +- 只允许任务原关联设备使用自身 Device Token 重置 `completed`、`completed_partial` 或 `failed` 采集任务;跨设备统一返回 `TASK_NOT_FOUND`。 +- `requestId` 必须为 UUID。同一请求重复提交不会再次清空结果,响应中的 `replayed` 为 `true`。 +- 服务端在事务内锁定设备和任务,重新校验设备在线、同商品无其他 `pending/running` 采集任务,并确认该设备没有正在执行或持有有效租约的其他采集/采购任务。 +- 成功后保留 URL、goods_id、设备和规则快照,事务删除旧结构化结果及错误,把原任务恢复为 `pending`;Agent 随后只能通过既有 `next → claim → start` 调度执行。 +- 响应只返回 `taskId`、`status` 和可选的 `replayed`,不返回规则快照、URL、Token、控件树或截图。 +- 设备离线、任务非终态、设备忙或同商品存在活动任务时返回明确冲突,不支持离线排队。 diff --git a/server/app/goauto/task/agent_reset_test.go b/server/app/goauto/task/agent_reset_test.go new file mode 100644 index 0000000..b66b576 --- /dev/null +++ b/server/app/goauto/task/agent_reset_test.go @@ -0,0 +1,124 @@ +package task + +import ( + "context" + "errors" + "testing" + + "go-admin/app/goauto/device" + "go-admin/app/goauto/models" + + "github.com/google/uuid" + "gorm.io/gorm" +) + +func TestAgentResetIsDeviceScopedIdleAndIdempotent(t *testing.T) { + db := openTaskDatabase(t) + deviceA, tokenA := registerTaskDevice(t, db, "device-a") + _, tokenB := registerTaskDevice(t, db, "device-b") + target := createTask(t, db, &deviceA.ID) + oldTitle := "旧结果" + if err := db.Session(&gorm.Session{SkipHooks: true}).Model(&models.CollectionTask{}).Where("id = ?", target.ID).Updates(map[string]any{ + "status": models.TaskStatusFailed, "active_slot": nil, "title": oldTitle, + "error_code": "RULE_NOT_MATCHED", "error_message": "没有找到规格入口", + }).Error; err != nil { + t.Fatal(err) + } + if err := db.Create(&models.CollectionColorPrice{TaskID: target.ID, Color: "黑色", PriceCent: 1000}).Error; err != nil { + t.Fatal(err) + } + service := newTaskService(db) + request := ActionRequest{RequestID: uuid.NewString()} + _, err := service.ResetForDevice(context.Background(), target.ID, request, "") + var deviceError *device.ServiceError + if !errors.As(err, &deviceError) || deviceError.Code != device.CodeTokenInvalid { + t.Fatalf("missing token was accepted: %v", err) + } + if _, err := service.ResetForDevice(context.Background(), target.ID, request, tokenB); taskErrorCode(t, err) != CodeTaskNotFound { + t.Fatalf("other device reset task: %v", err) + } + + busy := createTask(t, db, &deviceA.ID) + one := uint8(1) + if err := db.Session(&gorm.Session{SkipHooks: true}).Model(&models.CollectionTask{}).Where("id = ?", busy.ID).Updates(map[string]any{ + "status": models.TaskStatusRunning, "active_slot": one, "device_run_slot": one, + }).Error; err != nil { + t.Fatal(err) + } + if _, err := service.ResetForDevice(context.Background(), target.ID, request, tokenA); taskErrorCode(t, err) != CodeDeviceBusy { + t.Fatalf("busy device reset task: %v", err) + } + if err := db.Session(&gorm.Session{SkipHooks: true}).Model(&models.CollectionTask{}).Where("id = ?", busy.ID).Updates(map[string]any{ + "status": models.TaskStatusFailed, "active_slot": nil, "device_run_slot": nil, + }).Error; err != nil { + t.Fatal(err) + } + var product models.PDDProduct + if err := db.First(&product, target.PDDProductID).Error; err != nil { + t.Fatal(err) + } + purchaseBusy := models.PurchaseTask{ + PDDProductID: product.ID, DeviceID: &deviceA.ID, + ExecutionMode: models.PurchaseExecutionModeRehearsal, Status: models.PurchaseTaskStatusRunning, + PDDURLSnapshot: product.URL, PDDGoodsIDSnapshot: product.GoodsID, + Quantity: 1, ReferenceUnitPriceCent: 100, MinUnitPriceCent: 20, MaxUnitPriceCent: 150, + Currency: "CNY", RuleType: "pddPurchase", RuleSchemaVersion: 1, + RequiredCapabilitiesJSON: `[]`, RuleSnapshot: `{}`, CreateRequestID: uuid.NewString(), + } + if err := db.Create(&purchaseBusy).Error; err != nil { + t.Fatal(err) + } + if _, err := service.ResetForDevice(context.Background(), target.ID, request, tokenA); taskErrorCode(t, err) != CodeDeviceBusy { + t.Fatalf("running purchase did not block reset: %v", err) + } + if err := db.Session(&gorm.Session{SkipHooks: true}).Model(&models.PurchaseTask{}).Where("id = ?", purchaseBusy.ID).Updates(map[string]any{ + "status": models.PurchaseTaskStatusFailed, "active_slot": nil, "device_run_slot": nil, "account_run_slot": nil, + }).Error; err != nil { + t.Fatal(err) + } + + reset, err := service.ResetForDevice(context.Background(), target.ID, request, tokenA) + if err != nil || reset.TaskID != target.ID || reset.Status != models.TaskStatusPending || reset.Replayed { + t.Fatalf("reset: %+v %v", reset, err) + } + var updated models.CollectionTask + if err := db.First(&updated, target.ID).Error; err != nil { + t.Fatal(err) + } + if updated.DeviceID == nil || *updated.DeviceID != deviceA.ID || updated.Title != nil || updated.ErrorCode != nil { + t.Fatalf("reset did not preserve assignment or clear result: %+v", updated) + } + var prices int64 + if err := db.Model(&models.CollectionColorPrice{}).Where("task_id = ?", target.ID).Count(&prices).Error; err != nil || prices != 0 { + t.Fatalf("old prices remain: count=%d err=%v", prices, err) + } + replayed, err := service.ResetForDevice(context.Background(), target.ID, request, tokenA) + if err != nil || !replayed.Replayed { + t.Fatalf("reset replay: %+v %v", replayed, err) + } + if _, err := service.ResetForDevice(context.Background(), target.ID, ActionRequest{RequestID: uuid.NewString()}, tokenA); taskErrorCode(t, err) != CodeTaskStateConflict { + t.Fatalf("new request reset pending task: %v", err) + } +} + +func TestAgentResetRejectsActiveProductTask(t *testing.T) { + db := openTaskDatabase(t) + deviceRecord, token := registerTaskDevice(t, db, "device-one") + target := createTask(t, db, &deviceRecord.ID) + if err := db.Session(&gorm.Session{SkipHooks: true}).Model(&models.CollectionTask{}).Where("id = ?", target.ID).Updates(map[string]any{ + "status": models.TaskStatusCompleted, "active_slot": nil, + }).Error; err != nil { + t.Fatal(err) + } + conflict := models.CollectionTask{ + PDDProductID: target.PDDProductID, RuleID: target.RuleID, DeviceID: &deviceRecord.ID, + Status: models.TaskStatusPending, URLSnapshot: target.URLSnapshot, + GoodsIDSnapshot: target.GoodsIDSnapshot, RuleSnapshot: target.RuleSnapshot, + } + if err := db.Create(&conflict).Error; err != nil { + t.Fatal(err) + } + if _, err := newTaskService(db).ResetForDevice(context.Background(), target.ID, ActionRequest{RequestID: uuid.NewString()}, token); taskErrorCode(t, err) != CodeProductTaskActive { + t.Fatalf("active product task was ignored: %v", err) + } +} diff --git a/server/app/goauto/task/handler.go b/server/app/goauto/task/handler.go index b920fd0..e2421e4 100644 --- a/server/app/goauto/task/handler.go +++ b/server/app/goauto/task/handler.go @@ -80,6 +80,30 @@ func (handler Handler) AgentHistoryDetail(context *gin.Context) { context.JSON(http.StatusOK, gin.H{"data": result}) } +func (handler Handler) AgentReset(context *gin.Context) { + id, err := taskID(context) + if err != nil || id == 0 { + writeError(context, serviceError(device.CodeInvalidRequest, "taskId 无效")) + return + } + var request ActionRequest + if err := decodeStrict(context, &request); err != nil { + writeError(context, serviceError(device.CodeInvalidRequest, "请求 JSON 无效")) + return + } + service, token, ok := handler.service(context) + if !ok { + return + } + result, err := service.ResetForDevice(context.Request.Context(), id, request, token) + if err != nil { + writeError(context, err) + return + } + context.Header("Cache-Control", "no-store") + context.JSON(http.StatusOK, gin.H{"data": result}) +} + func (handler Handler) Claim(context *gin.Context) { handler.action(context, (*Service).Claim) } func (handler Handler) Start(context *gin.Context) { handler.action(context, (*Service).Start) } diff --git a/server/app/goauto/task/lifecycle_service.go b/server/app/goauto/task/lifecycle_service.go index 72e0a8b..d8297bb 100644 --- a/server/app/goauto/task/lifecycle_service.go +++ b/server/app/goauto/task/lifecycle_service.go @@ -3,7 +3,9 @@ package task import ( "context" "errors" + "time" + "go-admin/app/goauto/device" "go-admin/app/goauto/models" "gorm.io/gorm" @@ -16,12 +18,44 @@ type DeleteResponse struct { Replayed bool `json:"replayed,omitempty"` } +type AgentResetResponse struct { + TaskID uint64 `json:"taskId"` + Status string `json:"status"` + Replayed bool `json:"replayed,omitempty"` +} + func (service *Service) Reset(ctx context.Context, taskID uint64, request ActionRequest) (DetailResponse, error) { + return service.reset(ctx, taskID, request, nil) +} + +func (service *Service) ResetForDevice(ctx context.Context, taskID uint64, request ActionRequest, token string) (AgentResetResponse, error) { + detail, err := service.reset(ctx, taskID, request, &token) + if err != nil { + return AgentResetResponse{}, err + } + return AgentResetResponse{TaskID: detail.Task.ID, Status: detail.Task.Status, Replayed: detail.Replayed}, nil +} + +func (service *Service) reset(ctx context.Context, taskID uint64, request ActionRequest, deviceToken *string) (DetailResponse, error) { if err := validateAction(taskID, request); err != nil { return DetailResponse{}, err } var replayed bool err := service.DB.WithContext(ctx).Transaction(func(tx *gorm.DB) error { + var deviceRecord *models.AgentDevice + if deviceToken != nil { + authenticated, err := device.NewService(tx).Authenticate(ctx, *deviceToken) + if err != nil { + return err + } + if err := tx.Clauses(clause.Locking{Strength: "UPDATE"}).First(&authenticated, authenticated.ID).Error; err != nil { + return internalError(err) + } + if authenticated.Status != models.DeviceStatusOnline { + return serviceError(CodeDeviceOffline, "设备离线,不能重新采集") + } + deviceRecord = &authenticated + } var record models.CollectionTask if err := tx.Clauses(clause.Locking{Strength: "UPDATE"}).First(&record, taskID).Error; err != nil { if errors.Is(err, gorm.ErrRecordNotFound) { @@ -29,6 +63,9 @@ func (service *Service) Reset(ctx context.Context, taskID uint64, request Action } return internalError(err) } + if deviceRecord != nil && (record.DeviceID == nil || *record.DeviceID != deviceRecord.ID) { + return serviceError(CodeTaskNotFound, "采集任务不存在") + } if record.ResetRequestID != nil && *record.ResetRequestID == request.RequestID { replayed = true return nil @@ -47,6 +84,11 @@ func (service *Service) Reset(ctx context.Context, taskID uint64, request Action if active > 0 { return serviceError(CodeProductTaskActive, "该商品已有待执行或执行中的任务") } + if deviceRecord != nil { + if err := ensureDeviceIdleForReset(tx, deviceRecord.ID, record.ID, service.Now()); err != nil { + return err + } + } if err := deleteResultChildren(tx, record.ID); err != nil { return err } @@ -71,6 +113,29 @@ func (service *Service) Reset(ctx context.Context, taskID uint64, request Action return detail, err } +func ensureDeviceIdleForReset(tx *gorm.DB, deviceID, taskID uint64, now time.Time) error { + var busy int64 + if err := tx.Model(&models.CollectionTask{}). + Where("id <> ? AND device_id = ? AND (status = ? OR (status = ? AND lease_expires_at > ?))", taskID, deviceID, models.TaskStatusRunning, models.TaskStatusPending, now). + Count(&busy).Error; err != nil { + return internalError(err) + } + if busy > 0 { + return serviceError(CodeDeviceBusy, "设备正在执行其他任务") + } + if err := tx.Model(&models.PurchaseTask{}). + Where("device_id = ? AND (status IN ? OR (status IN ? AND lease_expires_at > ?))", deviceID, + []string{models.PurchaseTaskStatusRunning, models.PurchaseTaskStatusOrderSubmitStarted}, + []string{models.PurchaseTaskStatusPending, models.PurchaseTaskStatusSpecProbePending}, now). + Count(&busy).Error; err != nil { + return internalError(err) + } + if busy > 0 { + return serviceError(CodeDeviceBusy, "设备正在执行其他任务") + } + return nil +} + func (service *Service) DeleteFailed(ctx context.Context, taskID uint64, request ActionRequest) (DeleteResponse, error) { if err := validateAction(taskID, request); err != nil { return DeleteResponse{}, err diff --git a/server/app/goauto/task/router.go b/server/app/goauto/task/router.go index 8d4a654..b146075 100644 --- a/server/app/goauto/task/router.go +++ b/server/app/goauto/task/router.go @@ -19,6 +19,7 @@ func InitRouter(engine *gin.Engine, auth *jwt.GinJWTMiddleware) { agent.GET("/tasks/next", handler.Next) agent.GET("/collection-tasks", handler.AgentHistory) agent.GET("/collection-tasks/:taskId", handler.AgentHistoryDetail) + agent.POST("/collection-tasks/:taskId/reset", handler.AgentReset) agent.POST("/tasks/:taskId/claim", handler.Claim) agent.POST("/tasks/:taskId/start", handler.Start) agent.POST("/tasks/:taskId/result", handler.Result)