feat(#155): reuse failed collection tasks by attempt

This commit is contained in:
QiuSW
2026-08-29 16:05:37 +08:00
parent 78ab8716a6
commit cc326ebc6c
20 changed files with 520 additions and 118 deletions
+2 -2
View File
@@ -11,8 +11,8 @@ android {
applicationId = "cn.ilapage.goauto.agent"
minSdk = 23
targetSdk = 34
versionCode = 39
versionName = "0.9.26"
versionCode = 40
versionName = "0.9.27"
testInstrumentationRunner = "androidx.test.runner.AndroidJUnitRunner"
@@ -35,6 +35,7 @@ data class HeartbeatResult(
data class AgentTask(
val taskId: Long,
val attemptNumber: Int,
val pddProductId: Long?,
val urlSnapshot: String,
val goodsIdSnapshot: String,
@@ -78,6 +79,7 @@ data class PurchaseAgentTask(
data class CollectionHistoryItem(
val taskId: Long,
val attemptNumber: Int,
val status: String,
val source: String,
val goodsId: String,
@@ -95,6 +97,7 @@ data class HistoryColorPrice(val color: String, val priceCent: Long)
data class HistorySku(val specs: Map<String, String>, val priceCent: Long, val available: Boolean, val complete: Boolean)
data class CollectionHistoryDetail(
val task: CollectionHistoryItem,
val attempts: List<CollectionHistoryAttempt>,
val shopName: String?,
val salesText: String?,
val reviewCount: Long?,
@@ -109,7 +112,17 @@ data class CollectionHistoryDetail(
val replacementActivationErrorMessage: String?,
)
data class CollectionResetResult(val taskId: Long, val status: String, val replayed: Boolean)
data class CollectionHistoryAttempt(
val attemptNumber: Int,
val status: String,
val ruleId: Long,
val errorCode: String?,
val errorMessage: String?,
val finishedAt: String?,
val archivedAt: String,
)
data class CollectionResetResult(val taskId: Long, val attemptNumber: Int, val status: String, val replayed: Boolean)
data class PurchaseRetryResult(
val sourceTaskId: Long,
@@ -334,6 +347,14 @@ class AgentApiClient(private val serverUrl: String) {
val data = requireNotNull(request("GET", "/api/agent/v1/collection-tasks/$taskId", null, token)).getJSONObject("data")
return CollectionHistoryDetail(
task = collectionHistoryItem(data.getJSONObject("task")),
attempts = data.optJSONArray("attempts")?.objects { item ->
CollectionHistoryAttempt(
attemptNumber = item.optInt("attemptNumber", 1), status = item.getString("status"),
ruleId = item.getLong("ruleId"), errorCode = item.nullableString("errorCode"),
errorMessage = item.nullableString("errorMessage"), finishedAt = item.nullableString("finishedAt"),
archivedAt = item.getString("archivedAt"),
)
}.orEmpty(),
shopName = data.nullableString("shopName"), salesText = data.nullableString("salesText"),
reviewCount = data.nullableLong("reviewCount"),
dimensions = data.getJSONArray("dimensions").objects { item -> HistoryDimension(item.getString("key"), item.getString("name"), item.getJSONArray("values").strings()) },
@@ -355,7 +376,7 @@ 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"))
return CollectionResetResult(data.getLong("taskId"), data.optInt("attemptNumber", 1), data.getString("status"), data.optBoolean("replayed"))
}
fun purchaseHistory(token: String, page: Int, status: String?, taskNo: String?, days: Int = 30, pageSize: Int = 20): HistoryPage<PurchaseHistoryItem> {
@@ -398,7 +419,7 @@ class AgentApiClient(private val serverUrl: String) {
}
private fun collectionHistoryItem(data: JSONObject) = CollectionHistoryItem(
taskId = data.getLong("taskId"), status = data.getString("status"), source = data.optString("source", "admin"), goodsId = data.getString("goodsId"),
taskId = data.getLong("taskId"), attemptNumber = data.optInt("attemptNumber", 1), status = data.getString("status"), source = data.optString("source", "admin"), goodsId = data.getString("goodsId"),
title = data.nullableString("title"), missingCount = data.optInt("missingCount"),
errorCode = data.nullableString("errorCode"), errorMessage = data.nullableString("errorMessage"),
finishedAt = data.nullableString("finishedAt"), createdAt = data.getString("createdAt"),
@@ -419,6 +440,7 @@ class AgentApiClient(private val serverUrl: String) {
private fun task(data: JSONObject) = AgentTask(
taskId = data.getLong("taskId"),
attemptNumber = data.optInt("attemptNumber", 1),
pddProductId = data.nullableLong("pddProductId"),
urlSnapshot = data.getString("urlSnapshot"),
goodsIdSnapshot = data.getString("goodsIdSnapshot"),
@@ -13,7 +13,7 @@ class TaskHistoryCache(context: Context) {
fun saveCollection(days: Int, items: List<CollectionHistoryItem>) = save(COLLECTION, days, JSONArray().apply {
items.forEach { item -> put(JSONObject()
.put("taskId", item.taskId).put("status", item.status).put("source", item.source).put("goodsId", item.goodsId)
.put("taskId", item.taskId).put("attemptNumber", item.attemptNumber).put("status", item.status).put("source", item.source).put("goodsId", item.goodsId)
.putNullable("title", item.title).put("missingCount", item.missingCount)
.putNullable("errorCode", item.errorCode).putNullable("errorMessage", item.errorMessage)
.putNullable("finishedAt", item.finishedAt).put("createdAt", item.createdAt)) }
@@ -34,6 +34,7 @@ class TaskHistoryCache(context: Context) {
fun collection(days: Int, page: Int, status: String?, taskNo: String?): HistoryPage<CollectionHistoryItem>? {
val values = read(COLLECTION, days)?.map { data -> CollectionHistoryItem(
taskId = data.getLong("taskId"),
attemptNumber = data.optInt("attemptNumber", 1),
status = data.getString("status"),
source = data.optString("source", "admin"),
goodsId = data.optString("goodsId"),
@@ -66,15 +66,13 @@ internal object TaskSearchInputPolicy {
internal object CollectionResetPolicy {
private val terminalStatuses = setOf("completed", "completed_partial", "failed")
fun isAllowed(status: String, source: String): Boolean = status in terminalStatuses && source != "agent_current_page"
fun isAllowed(status: String, source: String): Boolean =
status == "failed" || (source != "agent_current_page" && status in terminalStatuses)
fun showsListAction(status: String, source: String): Boolean = status == "failed" && source != "agent_current_page"
fun showsListAction(status: String): Boolean = status == "failed"
fun confirmationMessage(status: String): String = if (status == "failed") {
"重新采集会清除本任务现有的错误和采集结果,然后重新加入当前设备的任务队列。"
} else {
"重新采集会清除本任务现有的全部采集结果,然后重新加入当前设备的任务队列。"
}
fun confirmationMessage(status: String): String =
"将复用原任务并按最新有效规则重新采集。当前${if (status == "failed") "错误和结果" else "采集结果"}会归档到执行记录,不会创建新任务。"
}
internal object PurchaseRetryPolicy {
@@ -411,10 +409,10 @@ class TaskHistoryFragment : Fragment() {
}
return context.card(context.cardColumn().apply {
val goodsLabel = item.goodsId.ifBlank { "未识别商品" }
addView(context.label("#${item.taskId} · $goodsLabel", 17f, context.getColor(R.color.agent_text), true))
addView(context.label("#${item.taskId} · 第 ${item.attemptNumber} 次 · $goodsLabel", 17f, context.getColor(R.color.agent_text), true))
addView(context.label("$status · ${formatTime(item.finishedAt ?: item.createdAt)}", 13f, statusColor(item.status)).apply { setPadding(0, context.dp(4), 0, 0) })
addView(context.label(summary, 14f, context.getColor(R.color.agent_text_muted)).apply { setPadding(0, context.dp(8), 0, 0) })
if (CollectionResetPolicy.showsListAction(item.status, item.source)) {
if (CollectionResetPolicy.showsListAction(item.status)) {
addView(LinearLayout(context).apply {
gravity = Gravity.END
addView(MaterialButton(context, null, com.google.android.material.R.attr.materialButtonOutlinedStyle).apply {
@@ -523,13 +521,20 @@ class TaskHistoryFragment : Fragment() {
append("$prices\n\n$dimensions\n\nSKU:${detail.skus.size} 条")
}
resultColumn.addView(context.card(context.cardColumn().apply {
addView(context.label("#${task.taskId} · ${task.goodsId.ifBlank { "未识别商品" }}", 20f, context.getColor(R.color.agent_text), true))
addView(context.label("#${task.taskId} · 第 ${task.attemptNumber} 次 · ${task.goodsId.ifBlank { "未识别商品" }}", 20f, context.getColor(R.color.agent_text), true))
if (task.source == "agent_current_page") addView(context.label("来源:Agent 当前页面", 13f, context.getColor(R.color.agent_text_muted)))
addView(context.label(collectionStatus(task.status), 14f, statusColor(task.status)).apply { setPadding(0, context.dp(4), 0, context.dp(12)) })
addView(context.label(info, 14f))
}), 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 (detail.attempts.isNotEmpty()) {
val history = detail.attempts.joinToString("\n") { attempt ->
val error = attempt.errorMessage?.let { " · $it" }.orEmpty()
"第 ${attempt.attemptNumber} 次 · ${collectionStatus(attempt.status)}$error"
}
resultColumn.addView(context.centeredMessage("历史执行记录", history))
}
if (CollectionResetPolicy.isAllowed(task.status, task.source)) {
resultColumn.addView(MaterialButton(context).apply {
text = "重新采集"
@@ -620,7 +625,8 @@ class TaskHistoryFragment : Fragment() {
}
private fun confirmReset(task: CollectionHistoryItem) {
collectionCooldownMessage()?.let { message ->
val problem = if (task.source == "agent_current_page") currentPageCollectionProblem() else collectionCooldownMessage()
problem?.let { message ->
MaterialAlertDialogBuilder(requireContext())
.setTitle("暂时不能重新采集")
.setMessage(message)
@@ -632,13 +638,15 @@ class TaskHistoryFragment : Fragment() {
.setTitle("重新采集 #${task.taskId}?")
.setMessage(CollectionResetPolicy.confirmationMessage(task.status))
.setNegativeButton("取消", null)
.setPositiveButton("确认重新采集") { _, _ -> resetCollectionTask(task.taskId) }
.setPositiveButton("确认重新采集") { _, _ -> resetCollectionTask(task) }
.show()
}
private fun resetCollectionTask(taskId: Long) {
private fun resetCollectionTask(task: CollectionHistoryItem) {
val taskId = task.taskId
val context = requireContext()
collectionCooldownMessage()?.let { message ->
val problem = if (task.source == "agent_current_page") currentPageCollectionProblem() else collectionCooldownMessage()
problem?.let { message ->
showMessage("暂时不能重新采集", message, "返回任务详情") { loadCollectionDetail(taskId) }
return
}
@@ -658,7 +666,7 @@ class TaskHistoryFragment : Fragment() {
AgentForegroundService.start(requireContext())
showMessage(
"已加入重新采集队列",
"任务 #${result.taskId} 将由当前设备通过正常调度重新执行。",
"任务 #${result.taskId} 将以第 ${result.attemptNumber} 次执行,由当前设备通过正常调度重新采集。",
"返回采集记录",
) { page = 1; load() }
}
@@ -13,22 +13,22 @@ class CollectionResetPolicyTest {
assertTrue(CollectionResetPolicy.isAllowed("failed", "admin"))
assertFalse(CollectionResetPolicy.isAllowed("pending", "admin"))
assertFalse(CollectionResetPolicy.isAllowed("running", "admin"))
assertFalse(CollectionResetPolicy.isAllowed("failed", "agent_current_page"))
assertTrue(CollectionResetPolicy.isAllowed("failed", "agent_current_page"))
assertFalse(CollectionResetPolicy.isAllowed("completed", "agent_current_page"))
}
@Test
fun `confirmation always explains result deletion`() {
assertTrue(CollectionResetPolicy.confirmationMessage("failed").contains("清除"))
assertTrue(CollectionResetPolicy.confirmationMessage("completed").contains("全部采集结果"))
assertTrue(CollectionResetPolicy.confirmationMessage("failed").contains("归档"))
assertTrue(CollectionResetPolicy.confirmationMessage("completed").contains("最新有效规则"))
}
@Test
fun `only failed collection tasks show the list reset action`() {
assertTrue(CollectionResetPolicy.showsListAction("failed", "admin"))
assertFalse(CollectionResetPolicy.showsListAction("completed", "admin"))
assertFalse(CollectionResetPolicy.showsListAction("completed_partial", "admin"))
assertFalse(CollectionResetPolicy.showsListAction("pending", "admin"))
assertFalse(CollectionResetPolicy.showsListAction("running", "admin"))
assertFalse(CollectionResetPolicy.showsListAction("failed", "agent_current_page"))
assertTrue(CollectionResetPolicy.showsListAction("failed"))
assertFalse(CollectionResetPolicy.showsListAction("completed"))
assertFalse(CollectionResetPolicy.showsListAction("completed_partial"))
assertFalse(CollectionResetPolicy.showsListAction("pending"))
assertFalse(CollectionResetPolicy.showsListAction("running"))
}
}
+9 -7
View File
@@ -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: 61a6effd4badcf94d9c4745521c4f3945a7be3ef
synchronized_at: 2026-08-29T03:46:32Z
wiki_revision: fdc3e2871913ea7c64fa9729296ce41273866c2b
synchronized_at: 2026-08-29T08:00:07Z
<!-- gitea-wiki-mirror:end -->
<!-- gitea-wiki-mirror:start -->
@@ -132,7 +132,7 @@ Android Portal/Agent
| 采购任务数据与类型化规则契约 | `server/app/goauto/models/purchase.go`、`server/app/goauto/purchasecontract/`;迁移 `server/cmd/migrate/migration/version-local/1786701100000_purchase_contract.go` |
| 采购任务单条/批量预检与创建、租约、attempt 幂等、Admin 只读查询和人工处置状态机 | `server/app/goauto/purchase/`;`batch-preview` 通过批量预加载 SYB、蝦皮、PDD 与最新任务执行快速只读预检,不访问 AI;`batch-create` 仍按最新数据逐项完整复核并可在必要时调用 AI;既有追加迁移为 `server/cmd/migrate/migration/version-local/1786701200000_purchase_state_machine.go` |
| 采购规格标准化匹配与 AI Provider 设置 | `server/app/goauto/aimatching/`;服务端先做繁简、空白/全半角/大小写和公斤/斤的唯一确定性匹配,再按需调用单一 OpenAI-compatible Provider;`1786701300000_ai_matching_setting.go` 创建设置表,`1786701400000_ai_matching_setting_plain_api_key.go` 将原加密列迁移为内部明文 `api_key`,仅管理员读取 |
| 任务领取、结果、重置与删除 | `server/app/goauto/task/` |
| 任务领取、结果、同任务 attempt 重采与删除 | `server/app/goauto/task/`;`collection_task.attempt_number` 保存当前次数,`collection_task_attempt` 归档旧终态执行,增量迁移为 `server/cmd/migrate/migration/version-local/1787983500000_collection_task_attempt.go` |
| 管理端基线 | `web/`(go-admin-ui v3.0.0) |
| 管理端闭环页面 | `web/src/views/goauto/` |
| 备货采购服务端路径 | `purchase_task.task_type` 与迁移 `1787885300000_stock_purchase.go`;`POST /api/admin/v1/purchase-tasks/stock` 由 `server/app/goauto/purchase/service.go` 校验 PDD 当前可选规格并固化 `direct_select`,复用既有任务状态机和设备/账号互斥;重试、替换和 SYB 回填显式排除 `stock` |
@@ -184,11 +184,13 @@ Android Portal/Agent
- 服务端采集与采购模块各自使用 `agent_history.go` 负责 Device Token 设备隔离、最近 30 天过滤、分页和最小化 DTO。
- `purchase_task.actual_unit_price_cent` 是可空字段,只记录 Agent 实际观察到的采购单价;旧任务不回填。
## Agent 受控重新采集(#91)
## Agent 受控重新采集(#91、#155)
- 服务端 `task/lifecycle_service.go` 的同一 Reset 事务同时服务管理端和设备端;设备端入口使用显式 Device Token 模式,不能用空 Token 退化为管理端路径。
- `POST /api/agent/v1/collection-tasks/{taskId}/reset` 只返回最小状态,设备归属、忙碌和商品冲突在事务内检查。
- Android `TaskHistoryFragment` 只在采集终态详情显示“重新采集”,确认成功后调用 `AgentForegroundService.start` 唤醒现有轮询,任务执行仍由 `next/claim/start` 完成。
- Reset 复用原 `collection_task.id`:事务先归档旧终态到 `collection_task_attempt`,递增 `attempt_number`,按来源刷新最新规则快照,再清理主表当前结果/错误/租约并恢复 `pending`。归档以 task + attempt 唯一,reset requestId 也唯一,阻止并发或重放重复递增。
- Admin 来源读取原规则当前存活内容;Agent 当前页来源读取当前手动默认规则。当前页未识别任务可以在下一 attempt 首次绑定,已识别任务仍由 `IdentifyCurrentPage` 的 goods_id 冲突保护限制为原商品。
- `POST /api/agent/v1/collection-tasks/{taskId}/reset` 返回 `taskId`、`attemptNumber`、状态和重放标记;设备归属/在线/忙碌、规则兼容、商品冲突和采购占用在事务内检查。
- Android `TaskHistoryFragment` 对所有失败来源显示“重新采集”;当前页来源复用无障碍、近期 PDD 前台、本地互斥和冷却检查,成功后调用 `AgentForegroundService.start`,执行仍由 `next/claim/start` 完成。Admin 详情抽屉与 Agent 详情都展示 attempt 历史摘要;Agent 接口不返回历史规则快照。
## Agent 受控采购重试(#95)
@@ -222,7 +224,7 @@ Android Portal/Agent
- 服务端入口位于 `server/app/goauto/task/current_page.go`:`CreateCurrentPage` 原子完成设备校验、任务创建、规则快照和租约;`IdentifyCurrentPage` 对 Agent 上报的白名单链接重新校验并最终裁决 goods_id,负责 PDD 商品复用/创建和任务身份绑定;服务端保留受限正文解析作为短链兜底。路由为 `/api/agent/v1/current-page-collection-tasks` 及其 `/{taskId}/identify`。
- 默认规则设置位于 `server/app/goauto/rule/agent_manual_setting.go`,使用单行表 `agent_manual_collection_setting`。迁移 `1787790000000_agent_current_page_collection.go` 增加任务来源、识别幂等字段、可空 PDD 外键和默认规则表,并对旧任务回填 `admin`。
- `collection_task.source` 目前只允许 `admin` / `agent_current_page`。当前页面任务创建时由设备运行槽防并发,识别商品后再参与 PDD 商品活动槽;结构化结果仍写入既有任务及规格子表,不新增临时任务表。
- `collection_task.source` 目前只允许 `admin` / `agent_current_page`。当前页面任务创建时由设备运行槽防并发,识别商品后再参与 PDD 商品活动槽;结构化结果仍写入既有任务及规格子表,不新增临时任务表。失败后通过 #155 的同任务 attempt 机制重新进入队列,不创建第二条临时任务。
- Android 入口仍由 `TaskHistoryFragment` 承载;点击确认时先启动 `AgentForegroundService`,不先把 Agent 退到后台。服务先用 `TaskExecutionMutex` 预占本地串行槽,服务端创建成功后把预占转为真实 taskId;前台已是 PDD 时不执行返回,前台仍是 Agent 且近期见过 PDD 时只尝试一次返回,否则使用不清理任务栈、不带深链参数的 PDD 启动 Intent 恢复现有任务,并在有界等待确认 PDD 前台后才识别身份和调用既有 `PddProductDetailCollector`。
- `CurrentPageIdentityRunner` 负责页面证据、唯一分享/复制入口和面板清理;`ClipboardRelayActivity` 使用独立、不可导出的短生命周期任务在前台读取新鲜剪贴板,完成后移除中转任务并露出原 PDD 页面。`PddShareLinkExpander` 仅在手机侧用无 Cookie、无项目凭据的移动端 GET 有界展开白名单短链,逐跳校验并最多读取 64KB 正文;Agent 不裁决 goods_id。原始剪贴板、链接与响应正文不进入日志或缓存。
- v2 规则新增可选 `currentPageIdentity`(分享/复制别名和三个有界超时),旧 v2 规则由 Android 使用安全默认值;新设备能力为 `collector.pdd.current-page-share.v1`。
+7 -6
View File
@@ -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: d88d5d5aeabdb66bb67a2ec4b32da912fbf0629f
synchronized_at: 2026-08-29T07:11:22Z
wiki_revision: d8444503e7d5828d05fc7cdce5c6f503bb08adbc
synchronized_at: 2026-08-29T08:00:23Z
<!-- gitea-wiki-mirror:end -->
<!-- gitea-wiki-mirror:start -->
@@ -88,7 +88,7 @@ synchronized_at: 2026-08-29T01:39:44Z
- 规则创建后立即可用,不存在草稿、发布、版本或停用流程。
- 规则允许编辑;新任务读取编辑后的当前内容。
- 创建任务时保存完整 `rule_snapshot`,以后编辑规则不影响该任务。
- 创建任务时保存完整 `rule_snapshot`;规则编辑不影响正在执行的当前 attempt。人工重新采集会归档旧 attempt,并按来源刷新规则:Admin 来源读取原 `rule_id` 的当前存活内容,Agent 当前页来源读取当前手动采集默认规则;规则已删除、不可用或设备能力不兼容时拒绝重采且不修改原任务。
- 删除采用软删除;删除后不能创建新任务,但已有任务继续执行自身快照。
- 商品详情步骤应声明精确 `activityName`,并与包名和唯一控件共同作为页面证据;进入 PDD 登录 Activity 必须返回 `PDD_LOGIN_REQUIRED`,不能提交采集成功。
- 唯一文字节点不可点击时,Agent 只可点击其最近的可点击父容器;不得改点兄弟节点或相似文字。
@@ -108,9 +108,10 @@ synchronized_at: 2026-08-29T01:39:44Z
- 一台设备同时只执行一个任务。
- 同一 PDD 商品不能同时存在多个 `pending` 或 `running` 任务。
- 手机离线时运行中任务失败;默认不自动重试、不自动换机。
- 失败后允许人工重置或删除后重新创建。
- `completed`、`completed_partial`、`failed` 可以重置;`running` 禁止重置和删除。
- 重置保留原 URL、goods_id、规则和设备快照,在一个事务中删除旧规格/SKU、清空结果和错误,并恢复为 `pending`。
- 所有 `failed` 任务,无论由 Admin 还是 Agent 当前页创建,都允许采购员人工“重新采集”;必须复用原 `collection_task.id`,不得为同一次重采新建任务。Admin 来源继续允许重置 `completed`、`completed_partial`,Agent 当前页来源只允许失败后重采;`running` 禁止重置和删除。
- `collection_task.attempt_number` 表示当前执行次数。重采事务先把旧终态执行归档到 `collection_task_attempt`,保留当次来源、设备、规则 ID/快照、终态、结构化错误、结果摘要与时间,再递增 attempt、清理主表当前错误/结果/租约并恢复为 `pending`;不保存原始控件树、截图、Token 或个人数据。
- 重采使用最新有效规则:Admin 来源刷新原规则当前内容,Agent 当前页来源刷新当前手动采集默认规则。全部资格、规则、设备、商品冲突和幂等检查必须在事务写入前通过;相同 `requestId` 重放不得重复归档或递增。
- Agent 当前页任务在身份识别前失败时,下一 attempt 允许首次绑定商品;已经绑定商品时必须重新识别为相同 goods_id,不同商品明确失败且不得覆盖原身份。当前页发起重采前仍要求无障碍就绪、最近 PDD 前台证据、设备在线空闲及本地冷却通过。
- 任务自身结果仍只在任务详情查看;任务完成后会按完整/部分完成规则更新 PDD 商品最新档案,商品页不替代任务结果审计。
- PDD 商品列表批量操作只选择当前页,最多 100 个商品;每个商品仍创建一个独立采集任务。
- 已停用、没有可用采集规则或已有 `pending/running` 任务的商品不能勾选,并显示普通人可理解的原因。
+12 -11
View File
@@ -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: 0fbfd430daab3d703649dfab34100be107581b81
synchronized_at: 2026-08-29T07:11:55Z
wiki_revision: ec45e8baa2f55221bee3f9d6951379178020de17
synchronized_at: 2026-08-29T08:01:18Z
<!-- gitea-wiki-mirror:end -->
<!-- gitea-wiki-mirror:start -->
@@ -559,7 +559,7 @@ GET /api/agent/v1/purchase-tasks/{taskId}
- `pageSize` 最大为 50;`status` 与 `taskNo` 可以组合过滤。
- 采集任务编号允许 `35` 或 `#35`,采购任务编号允许 `12` 或不区分大小写的 `CG-12`;服务端按精确编号匹配。
- 列表和详情不返回 Device Token、URL、收货地址、规则快照或原始控件树。采购项额外返回服务端计算的 `retryable` 和可选 `retryDisabledReason`;除受控失败重试外,不提供取消、修改既有订单或支付入口。
- 采集详情返回标题、店铺、销量、评价数、规格维度、颜色价格、SKU、缺失项和结构化错误。
- 采集摘要返回当前 `attemptNumber`;采集详情返回标题、店铺、销量、评价数、规格维度、颜色价格、SKU、缺失项、结构化错误,以及不含规则快照的历史 attempt 序号、状态、规则 ID、错误和时间摘要。
- 采购详情返回蝦皮订单号、PDD 商品、已匹配颜色/尺码、数量、实际单价、PDD 订单号、下单时间和结构化错误。
- Agent 提交采购结果时可携带 `actualUnitPriceCent`(人民币分,非负)。服务端只保存 Agent 实际观察到的值;历史任务或未观察到价格的结果保持 `null`,客户端显示“未记录”。
@@ -573,13 +573,14 @@ Content-Type: application/json
{"requestId":"<uuid>"}
```
- 只允许任务原关联设备使用自身 Device Token 重置 `completed`、`completed_partial` 或 `failed` 采集任务;跨设备统一返回 `TASK_NOT_FOUND`。
- `requestId` 必须为 UUID。同一请求重复提交不会再次清空结果,响应中的 `replayed` 为 `true`。
- 只允许任务原关联设备使用自身 Device Token 重采:所有来源的 `failed` 任务均可重采;Admin 来源继续允许 `completed`、`completed_partial` 重置,Agent 当前页来源只允许失败后重采。跨设备统一返回 `TASK_NOT_FOUND`。
- `requestId` 必须为 UUID。同一请求重复提交不会重复归档、递增 attempt 或清空结果,响应中的 `replayed` 为 `true`。
- 服务端在事务内锁定设备和任务,重新校验设备在线、同商品无其他 `pending/running` 采集任务,并确认该设备没有正在执行或持有有效租约的其他采集/采购任务。
- 成功后保留 URL、goods_id、设备和规则快照,事务删除旧结构化结果及错误,把原任务恢复为 `pending`;Agent 随后只能通过既有 `next → claim → start` 调度执行。
- 响应只返回 `taskId`、`status` 和可选的 `replayed`,不返回规则快照、URL、Token、控件树或截图。
- 设备离线、任务非终态、设备忙或同商品存在活动任务时返回明确冲突,不支持离线排队。
- 成功时复用原 `collection_task.id`,把旧终态执行归档到 `collection_task_attempt`,将 `attemptNumber` 加一,再清理主表当前结果/错误/租约并恢复为 `pending`;不会新建采集任务。
- 新 attempt 使用最新有效规则:Admin 来源读取原 `ruleId` 的当前存活内容,Agent 当前页来源读取当前手动采集默认规则。规则不存在、不可用或与设备能力不兼容时拒绝且不修改任务。
- 当前页任务尚未识别商品时允许下一 attempt 首次绑定;已有商品身份时重新识别必须逐字命中原 goods_id,不同商品返回身份冲突。
- 响应返回 `taskId`、`attemptNumber`、`status` 和可选的 `replayed`,不返回规则快照、URL、Token、控件树或截图。
- 设备离线、任务非终态、设备忙、规则不可用或同商品存在活动任务时返回明确冲突,不支持离线排队。
## Agent 受控采购重试(#95)
@@ -626,7 +627,7 @@ Content-Type: application/json
- 只允许管理员读取和修改“Agent 手动采集默认规则”;PUT 的 `requestId` 必须为 UUID,重复请求幂等。
- 规则必须是仍存在的 PDD 商品详情 v2 规则。数据库迁移首次执行时,优先选择最近成功/部分成功任务使用的存活规则;没有成功历史时选择最近更新的存活 v2 规则;仍没有时保持未配置。
- 创建任务时重新验证规则与设备能力,并把完整规则快照固化到任务;以后修改默认规则不影响已创建任务。
- 创建任务时重新验证规则与设备能力,并把完整规则快照固化到第 1 次 attempt;以后修改默认规则不影响正在执行的 attempt,但失败任务人工重采时会读取当前默认规则并固化到新的 attempt。
### 创建并占用当前设备
@@ -658,7 +659,7 @@ Content-Type: application/json
- 服务端提取 5~32 位纯数字 `goods_id`,形成标准 URL,并在事务中创建或复用唯一 PDD 商品;同商品已有活动采集任务或身份冲突时拒绝。
- 相同识别 `requestId` 直接返回已确认身份且不再次访问短链;任务首次确认身份后,不允许不同请求覆盖为其他商品。
- 识别后继续复用 `POST /api/agent/v1/tasks/{taskId}/result` 与 `/fail`。结果接口要求任务已绑定身份且结果 goods_id 一致;完成、部分完成和失败继续按既有状态机释放设备槽。
- Agent 历史的任务摘要增加 `source`;`agent_current_page` 展示为“Agent 当前页面”,不新增任务状态。该来源任务不支持重置,失败后由用户从当前商品页重新发起新任务。
- Agent 历史的任务摘要返回 `source` 和当前 `attemptNumber`;详情返回不含规则快照的历史 attempt 摘要。所有 `failed` 当前页任务均可通过既有 reset 接口复用原任务重采;发起前 Android 必须确认无障碍就绪、没有本地执行中任务、最近见过 PDD 前台并通过采集冷却。
### 失败任务采集替代商品(#130)
+1
View File
@@ -52,6 +52,7 @@ func MigratedModels() []any {
&models.CollectionRule{},
&models.AgentManualCollectionSetting{},
&models.CollectionTask{},
&models.CollectionTaskAttempt{},
&models.CollectionDimension{},
&models.CollectionDimensionValue{},
&models.CollectionColorPrice{},
+30
View File
@@ -131,6 +131,7 @@ type CollectionTask struct {
URLSnapshot string `json:"urlSnapshot" gorm:"type:text;not null"`
GoodsIDSnapshot string `json:"goodsIdSnapshot" gorm:"size:32;not null;index"`
RuleSnapshot string `json:"ruleSnapshot" gorm:"type:text;not null"`
AttemptNumber int `json:"attemptNumber" gorm:"not null;default:1;check:ck_collection_task_attempt_number,attempt_number >= 1"`
LeaseExpiresAt *time.Time `json:"leaseExpiresAt" gorm:"index"`
LeaseVersion uint64 `json:"leaseVersion" gorm:"not null;default:0"`
ClaimRequestID *string `json:"-" gorm:"size:36;uniqueIndex:ux_collection_task_claim_request_id"`
@@ -170,6 +171,9 @@ func (task *CollectionTask) BeforeCreate(_ *gorm.DB) error {
if task.Source == "" {
task.Source = CollectionTaskSourceAdmin
}
if task.AttemptNumber == 0 {
task.AttemptNumber = 1
}
return task.syncGuardSlots()
}
@@ -223,6 +227,32 @@ func (task *CollectionTask) syncGuardSlots() error {
return nil
}
// CollectionTaskAttempt archives one terminal execution before the same
// collection task is reset. It deliberately stores only structured task
// facts; raw accessibility trees, screenshots and device secrets are never
// persisted here.
type CollectionTaskAttempt struct {
ID uint64 `json:"id" gorm:"primaryKey;autoIncrement"`
TaskID uint64 `json:"taskId" gorm:"not null;index;uniqueIndex:ux_collection_attempt_number,priority:1"`
Task CollectionTask `json:"-" gorm:"constraint:OnUpdate:CASCADE,OnDelete:CASCADE"`
AttemptNumber int `json:"attemptNumber" gorm:"not null;uniqueIndex:ux_collection_attempt_number,priority:2;check:ck_collection_attempt_number,attempt_number >= 1"`
Source string `json:"source" gorm:"size:32;not null"`
DeviceID *uint64 `json:"deviceId,omitempty" gorm:"index"`
RuleID uint64 `json:"ruleId" gorm:"not null;index"`
RuleSnapshot string `json:"ruleSnapshot" gorm:"type:text;not null"`
Status string `json:"status" gorm:"size:24;not null;index"`
ErrorCode *string `json:"errorCode,omitempty" gorm:"size:64"`
ErrorMessage *string `json:"errorMessage,omitempty" gorm:"size:1000"`
ResultSummaryJSON string `json:"-" gorm:"type:text;not null"`
StartedAt *time.Time `json:"startedAt,omitempty"`
FinishedAt *time.Time `json:"finishedAt,omitempty"`
ArchivedByResetRequestID string `json:"-" gorm:"size:36;not null;uniqueIndex:ux_collection_attempt_reset_request"`
ArchivedAt time.Time `json:"archivedAt" gorm:"not null"`
CreatedAt time.Time `json:"createdAt"`
}
func (CollectionTaskAttempt) TableName() string { return "collection_task_attempt" }
type CollectionDimension struct {
ID uint64 `json:"id" gorm:"primaryKey;autoIncrement"`
TaskID uint64 `json:"taskId" gorm:"not null;index;uniqueIndex:ux_collection_dimension_task_key,priority:1"`
+15 -9
View File
@@ -157,7 +157,7 @@ func TestSubmitResultPersistsChildrenAndIsIdempotent(t *testing.T) {
}
}
func TestResetClearsResultAndDeletedRuleStillWorks(t *testing.T) {
func TestResetRejectsDeletedRuleWithoutClearingResult(t *testing.T) {
db := openTaskDatabase(t)
deviceRecord, _ := registerTaskDevice(t, db, "device-one")
record := createTask(t, db, &deviceRecord.ID)
@@ -173,16 +173,22 @@ func TestResetClearsResultAndDeletedRuleStillWorks(t *testing.T) {
}
service := newTaskService(db)
request := ActionRequest{RequestID: uuid.NewString()}
detail, err := service.Reset(context.Background(), record.ID, request)
if err != nil {
t.Fatalf("reset: %v", err)
if _, err := service.Reset(context.Background(), record.ID, request); taskErrorCode(t, err) != CodeRuleNotFound {
t.Fatalf("deleted rule reset: %v", err)
}
if detail.Task.Status != models.TaskStatusPending || detail.Task.Title != nil || detail.Task.ErrorCode != nil || len(detail.ColorPrices) != 0 {
t.Fatalf("old result remains: %+v", detail)
var persisted models.CollectionTask
if err := db.First(&persisted, record.ID).Error; err != nil {
t.Fatal(err)
}
replay, err := service.Reset(context.Background(), record.ID, request)
if err != nil || !replay.Replayed {
t.Fatalf("bad reset replay: %+v %v", replay, err)
var prices, attempts int64
if err := db.Model(&models.CollectionColorPrice{}).Where("task_id = ?", record.ID).Count(&prices).Error; err != nil {
t.Fatal(err)
}
if err := db.Model(&models.CollectionTaskAttempt{}).Where("task_id = ?", record.ID).Count(&attempts).Error; err != nil {
t.Fatal(err)
}
if persisted.Status != models.TaskStatusFailed || persisted.ErrorCode == nil || prices != 1 || attempts != 0 {
t.Fatalf("rejected reset mutated task: task=%+v prices=%d attempts=%d", persisted, prices, attempts)
}
}
+47 -26
View File
@@ -26,17 +26,18 @@ type AgentHistoryRequest struct {
}
type AgentCollectionItem struct {
TaskID uint64 `json:"taskId"`
Status string `json:"status"`
Source string `json:"source"`
GoodsID string `json:"goodsId"`
Title *string `json:"title,omitempty"`
MissingCount int `json:"missingCount"`
ErrorCode *string `json:"errorCode,omitempty"`
ErrorMessage *string `json:"errorMessage,omitempty"`
StartedAt *time.Time `json:"startedAt,omitempty"`
FinishedAt *time.Time `json:"finishedAt,omitempty"`
CreatedAt time.Time `json:"createdAt"`
TaskID uint64 `json:"taskId"`
AttemptNumber int `json:"attemptNumber"`
Status string `json:"status"`
Source string `json:"source"`
GoodsID string `json:"goodsId"`
Title *string `json:"title,omitempty"`
MissingCount int `json:"missingCount"`
ErrorCode *string `json:"errorCode,omitempty"`
ErrorMessage *string `json:"errorMessage,omitempty"`
StartedAt *time.Time `json:"startedAt,omitempty"`
FinishedAt *time.Time `json:"finishedAt,omitempty"`
CreatedAt time.Time `json:"createdAt"`
}
type AgentCollectionList struct {
@@ -47,19 +48,31 @@ type AgentCollectionList struct {
}
type AgentCollectionDetail struct {
Task AgentCollectionItem `json:"task"`
ShopName *string `json:"shopName,omitempty"`
SalesText *string `json:"salesText,omitempty"`
ReviewCount *int64 `json:"reviewCount,omitempty"`
Dimensions []DetailDimension `json:"dimensions"`
ColorPrices []AgentColorPrice `json:"colorPrices"`
SKUs []AgentCollectionSKU `json:"skus"`
Missing []string `json:"missing"`
ReplacementEligible bool `json:"replacementEligible"`
ReplacementDisabledReason string `json:"replacementDisabledReason,omitempty"`
ReplacementMappingStatus string `json:"replacementMappingStatus,omitempty"`
ReplacementActivationStatus string `json:"replacementActivationStatus,omitempty"`
ReplacementActivationErrorMessage string `json:"replacementActivationErrorMessage,omitempty"`
Task AgentCollectionItem `json:"task"`
Attempts []AgentCollectionAttempt `json:"attempts"`
ShopName *string `json:"shopName,omitempty"`
SalesText *string `json:"salesText,omitempty"`
ReviewCount *int64 `json:"reviewCount,omitempty"`
Dimensions []DetailDimension `json:"dimensions"`
ColorPrices []AgentColorPrice `json:"colorPrices"`
SKUs []AgentCollectionSKU `json:"skus"`
Missing []string `json:"missing"`
ReplacementEligible bool `json:"replacementEligible"`
ReplacementDisabledReason string `json:"replacementDisabledReason,omitempty"`
ReplacementMappingStatus string `json:"replacementMappingStatus,omitempty"`
ReplacementActivationStatus string `json:"replacementActivationStatus,omitempty"`
ReplacementActivationErrorMessage string `json:"replacementActivationErrorMessage,omitempty"`
}
type AgentCollectionAttempt struct {
AttemptNumber int `json:"attemptNumber"`
Status string `json:"status"`
RuleID uint64 `json:"ruleId"`
ErrorCode *string `json:"errorCode,omitempty"`
ErrorMessage *string `json:"errorMessage,omitempty"`
StartedAt *time.Time `json:"startedAt,omitempty"`
FinishedAt *time.Time `json:"finishedAt,omitempty"`
ArchivedAt time.Time `json:"archivedAt"`
}
type AgentColorPrice struct {
@@ -155,8 +168,16 @@ func (service *Service) AgentHistoryDetail(ctx context.Context, taskID uint64, t
for _, value := range detail.SKUs {
skus = append(skus, AgentCollectionSKU{Specs: value.Specs, PriceCent: value.PriceCent, Available: value.Available, Complete: value.Complete})
}
attempts := make([]AgentCollectionAttempt, 0, len(detail.Attempts))
for _, value := range detail.Attempts {
attempts = append(attempts, AgentCollectionAttempt{
AttemptNumber: value.AttemptNumber, Status: value.Status, RuleID: value.RuleID,
ErrorCode: value.ErrorCode, ErrorMessage: value.ErrorMessage,
StartedAt: value.StartedAt, FinishedAt: value.FinishedAt, ArchivedAt: value.ArchivedAt,
})
}
return AgentCollectionDetail{
Task: agentCollectionItem(record), ShopName: record.ShopName, SalesText: record.SalesText,
Task: agentCollectionItem(record), Attempts: attempts, ShopName: record.ShopName, SalesText: record.SalesText,
ReviewCount: record.ReviewCount, Dimensions: detail.Dimensions, ColorPrices: colorPrices,
SKUs: skus, Missing: detail.Missing,
ReplacementEligible: inspection.Eligible, ReplacementDisabledReason: inspection.DisabledReason,
@@ -173,7 +194,7 @@ func agentCollectionItem(record models.CollectionTask) AgentCollectionItem {
missing = len(values)
}
return AgentCollectionItem{
TaskID: record.ID, Status: record.Status, Source: record.Source, GoodsID: record.GoodsIDSnapshot, Title: record.Title,
TaskID: record.ID, AttemptNumber: record.AttemptNumber, Status: record.Status, Source: record.Source, GoodsID: record.GoodsIDSnapshot, Title: record.Title,
MissingCount: missing, ErrorCode: record.ErrorCode, ErrorMessage: record.ErrorMessage,
StartedAt: record.StartedAt, FinishedAt: record.FinishedAt, CreatedAt: record.CreatedAt,
}
+70 -1
View File
@@ -4,6 +4,7 @@ import (
"context"
"errors"
"testing"
"time"
"go-admin/app/goauto/device"
"go-admin/app/goauto/models"
@@ -28,6 +29,14 @@ func TestAgentResetIsDeviceScopedIdleAndIdempotent(t *testing.T) {
t.Fatal(err)
}
service := newTaskService(db)
var originalRule models.CollectionRule
if err := db.First(&originalRule, target.RuleID).Error; err != nil {
t.Fatal(err)
}
const updatedRule = `{"steps":[{"action":"updated"}]}`
if err := db.Model(&originalRule).Update("content_json", updatedRule).Error; err != nil {
t.Fatal(err)
}
request := ActionRequest{RequestID: uuid.NewString()}
_, err := service.ResetForDevice(context.Background(), target.ID, request, "")
var deviceError *device.ServiceError
@@ -85,9 +94,16 @@ func TestAgentResetIsDeviceScopedIdleAndIdempotent(t *testing.T) {
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 {
if updated.DeviceID == nil || *updated.DeviceID != deviceA.ID || updated.Title != nil || updated.ErrorCode != nil || updated.AttemptNumber != 2 || updated.RuleSnapshot != updatedRule {
t.Fatalf("reset did not preserve assignment or clear result: %+v", updated)
}
var attempts []models.CollectionTaskAttempt
if err := db.Where("task_id = ?", target.ID).Find(&attempts).Error; err != nil || len(attempts) != 1 {
t.Fatalf("attempt archive: count=%d err=%v", len(attempts), err)
}
if attempts[0].AttemptNumber != 1 || attempts[0].Status != models.TaskStatusFailed || attempts[0].RuleSnapshot != target.RuleSnapshot || attempts[0].ErrorCode == nil || *attempts[0].ErrorCode != "RULE_NOT_MATCHED" {
t.Fatalf("unexpected attempt archive: %+v", attempts[0])
}
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)
@@ -96,11 +112,64 @@ func TestAgentResetIsDeviceScopedIdleAndIdempotent(t *testing.T) {
if err != nil || !replayed.Replayed {
t.Fatalf("reset replay: %+v %v", replayed, err)
}
if err := db.Where("task_id = ?", target.ID).Find(&attempts).Error; err != nil || len(attempts) != 1 {
t.Fatalf("replay duplicated attempt archive: count=%d err=%v", len(attempts), 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 TestAgentResetCurrentPageUnidentifiedTaskUsesLatestManualRule(t *testing.T) {
db := openTaskDatabase(t)
deviceRecord, token := registerTaskDeviceWithCapabilities(t, db, "current-page-reset", []string{
"rule.schema.v2", "collector.pdd.product-detail.v1", "collector.pdd.current-page-share.v1",
})
oldRule := models.CollectionRule{Name: "old", ContentJSON: v2TaskRuleSnapshot()}
latestRule := models.CollectionRule{Name: "latest", ContentJSON: v2TaskRuleSnapshot()}
if err := db.Create(&oldRule).Error; err != nil {
t.Fatal(err)
}
if err := db.Create(&latestRule).Error; err != nil {
t.Fatal(err)
}
if err := db.Create(&models.AgentManualCollectionSetting{ID: 1, RuleID: latestRule.ID}).Error; err != nil {
t.Fatal(err)
}
errorCode := "PDD_COPY_LINK_UNAVAILABLE"
errorMessage := "复制商品链接失败"
now := time.Now().UTC()
task := models.CollectionTask{
RuleID: oldRule.ID, DeviceID: &deviceRecord.ID, Source: models.CollectionTaskSourceAgentCurrentPage,
Status: models.TaskStatusFailed, URLSnapshot: "", GoodsIDSnapshot: "", RuleSnapshot: oldRule.ContentJSON,
ErrorCode: &errorCode, ErrorMessage: &errorMessage, FinishedAt: &now,
}
if err := db.Create(&task).Error; err != nil {
t.Fatal(err)
}
reset, err := newTaskService(db).ResetForDevice(context.Background(), task.ID, ActionRequest{RequestID: uuid.NewString()}, token)
if err != nil {
t.Fatalf("reset unidentified current-page task: %v", err)
}
if reset.TaskID != task.ID || reset.AttemptNumber != 2 || reset.Status != models.TaskStatusPending {
t.Fatalf("unexpected reset response: %+v", reset)
}
var updated models.CollectionTask
if err := db.First(&updated, task.ID).Error; err != nil {
t.Fatal(err)
}
if updated.PDDProductID != nil || updated.RuleID != latestRule.ID || updated.RuleSnapshot != latestRule.ContentJSON || updated.AttemptNumber != 2 {
t.Fatalf("current-page reset did not use latest rule or preserve empty identity: %+v", updated)
}
var attempt models.CollectionTaskAttempt
if err := db.Where("task_id = ? AND attempt_number = 1", task.ID).First(&attempt).Error; err != nil {
t.Fatal(err)
}
if attempt.ErrorCode == nil || *attempt.ErrorCode != errorCode || attempt.RuleID != oldRule.ID {
t.Fatalf("unexpected archived current-page attempt: %+v", attempt)
}
}
func TestAgentResetRejectsActiveProductTask(t *testing.T) {
db := openTaskDatabase(t)
deviceRecord, token := registerTaskDevice(t, db, "device-one")
+6
View File
@@ -254,6 +254,12 @@ func (service *Service) IdentifyCurrentPage(ctx context.Context, taskID uint64,
if record.GoodsIDSnapshot != resolved.GoodsID {
return serviceError(CodeCurrentPageIdentityConflict, "当前商品与任务已识别商品不一致")
}
now := service.Now()
if err := tx.Model(&models.CollectionTask{}).Where("id = ?", record.ID).Updates(map[string]any{
"identify_request_id": request.RequestID, "identity_resolved_at": now,
}).Error; err != nil {
return internalError(err)
}
response = CurrentPageIdentifyResponse{TaskID: record.ID, PDDProductID: *record.PDDProductID, GoodsID: record.GoodsIDSnapshot, URL: record.URLSnapshot, Replayed: true}
return nil
}
+134 -18
View File
@@ -2,6 +2,7 @@ package task
import (
"context"
"encoding/json"
"errors"
"time"
@@ -19,9 +20,10 @@ type DeleteResponse struct {
}
type AgentResetResponse struct {
TaskID uint64 `json:"taskId"`
Status string `json:"status"`
Replayed bool `json:"replayed,omitempty"`
TaskID uint64 `json:"taskId"`
AttemptNumber int `json:"attemptNumber"`
Status string `json:"status"`
Replayed bool `json:"replayed,omitempty"`
}
func (service *Service) Reset(ctx context.Context, taskID uint64, request ActionRequest) (DetailResponse, error) {
@@ -33,7 +35,7 @@ func (service *Service) ResetForDevice(ctx context.Context, taskID uint64, reque
if err != nil {
return AgentResetResponse{}, err
}
return AgentResetResponse{TaskID: detail.Task.ID, Status: detail.Task.Status, Replayed: detail.Replayed}, nil
return AgentResetResponse{TaskID: detail.Task.ID, AttemptNumber: detail.Task.AttemptNumber, Status: detail.Task.Status, Replayed: detail.Replayed}, nil
}
func (service *Service) reset(ctx context.Context, taskID uint64, request ActionRequest, deviceToken *string) (DetailResponse, error) {
@@ -66,45 +68,78 @@ func (service *Service) reset(ctx context.Context, taskID uint64, request Action
if deviceRecord != nil && (record.DeviceID == nil || *record.DeviceID != deviceRecord.ID) {
return serviceError(CodeTaskNotFound, "采集任务不存在")
}
if deviceRecord == nil && record.DeviceID != nil {
var assigned models.AgentDevice
if err := tx.Clauses(clause.Locking{Strength: "UPDATE"}).First(&assigned, *record.DeviceID).Error; err != nil {
return internalError(err)
}
if assigned.Status != models.DeviceStatusOnline {
return serviceError(CodeDeviceOffline, "设备离线,不能重新采集")
}
deviceRecord = &assigned
}
if record.ResetRequestID != nil && *record.ResetRequestID == request.RequestID {
replayed = true
return nil
}
if record.Source == models.CollectionTaskSourceAgentCurrentPage {
return serviceError(CodeTaskStateConflict, "当前页面采集不能重置,请回到商品页重新发起采集")
}
if record.PDDProductID == nil {
return serviceError(CodeTaskStateConflict, "任务尚未识别商品,不能重置")
var archivedReplay models.CollectionTaskAttempt
if err := tx.Where("archived_by_reset_request_id = ?", request.RequestID).First(&archivedReplay).Error; err == nil {
if archivedReplay.TaskID != record.ID {
return serviceError(CodeTaskStateConflict, "requestId 已被其他任务使用")
}
replayed = true
return nil
} else if !errors.Is(err, gorm.ErrRecordNotFound) {
return internalError(err)
}
if record.Status == models.TaskStatusRunning {
return serviceError(CodeTaskStateConflict, "执行中的任务不能重置")
}
if record.Source == models.CollectionTaskSourceAgentCurrentPage && record.Status != models.TaskStatusFailed {
return serviceError(CodeTaskStateConflict, "当前页面采集只允许失败任务重新采集")
}
if record.Status != models.TaskStatusCompleted && record.Status != models.TaskStatusCompletedPartial && record.Status != models.TaskStatusFailed {
return serviceError(CodeTaskStateConflict, "只有终态任务可以重置")
}
var active int64
if err := tx.Model(&models.CollectionTask{}).Where("id <> ? AND pdd_product_id = ? AND status IN ?", record.ID, record.PDDProductID,
[]string{models.TaskStatusPending, models.TaskStatusRunning}).Count(&active).Error; err != nil {
return internalError(err)
}
if active > 0 {
return serviceError(CodeProductTaskActive, "该商品已有待执行或执行中的任务")
if record.PDDProductID != nil {
var active int64
if err := tx.Model(&models.CollectionTask{}).Where("id <> ? AND pdd_product_id = ? AND status IN ?", record.ID, record.PDDProductID,
[]string{models.TaskStatusPending, models.TaskStatusRunning}).Count(&active).Error; err != nil {
return internalError(err)
}
if active > 0 {
return serviceError(CodeProductTaskActive, "该商品已有待执行或执行中的任务")
}
}
if deviceRecord != nil {
if err := ensureDeviceIdleForReset(tx, deviceRecord.ID, record.ID, service.Now()); err != nil {
return err
}
}
rule, err := currentResetRule(tx, record)
if err != nil {
return err
}
if deviceRecord != nil {
if err := ensureRuleCompatible(*deviceRecord, rule.ContentJSON); err != nil {
return err
}
}
now := service.Now()
if err := archiveCollectionAttempt(tx, record, request.RequestID, now); err != nil {
return err
}
if err := deleteResultChildren(tx, record.ID); err != nil {
return err
}
one := uint8(1)
updates := map[string]any{
"status": models.TaskStatusPending, "active_slot": one, "device_run_slot": nil,
"attempt_number": record.AttemptNumber + 1, "rule_id": rule.ID, "rule_snapshot": rule.ContentJSON,
"lease_expires_at": nil, "claim_request_id": nil, "start_request_id": nil,
"result_request_id": nil, "fail_request_id": nil, "reset_request_id": request.RequestID,
"result_request_id": nil, "fail_request_id": nil, "identify_request_id": nil, "reset_request_id": request.RequestID,
"title": nil, "shop_name": nil, "sales_text": nil, "review_count": nil, "missing_json": nil,
"error_code": nil, "error_message": nil, "started_at": nil, "finished_at": nil,
"error_code": nil, "error_message": nil, "started_at": nil, "finished_at": nil, "identity_resolved_at": nil,
}
if err := tx.Session(&gorm.Session{SkipHooks: true}).Model(&models.CollectionTask{}).Where("id = ?", record.ID).Updates(updates).Error; err != nil {
return internalError(err)
@@ -119,6 +154,87 @@ func (service *Service) reset(ctx context.Context, taskID uint64, request Action
return detail, err
}
func currentResetRule(tx *gorm.DB, record models.CollectionTask) (models.CollectionRule, error) {
var rule models.CollectionRule
if record.Source == models.CollectionTaskSourceAgentCurrentPage {
var setting models.AgentManualCollectionSetting
if err := tx.First(&setting, 1).Error; err != nil {
if errors.Is(err, gorm.ErrRecordNotFound) {
return rule, serviceError(CodeAgentManualRuleNotConfigured, "请先配置 Agent 手动采集规则")
}
return rule, internalError(err)
}
if err := tx.First(&rule, setting.RuleID).Error; err != nil {
if errors.Is(err, gorm.ErrRecordNotFound) {
return rule, serviceError(CodeAgentManualRuleNotConfigured, "请先配置 Agent 手动采集规则")
}
return rule, internalError(err)
}
if err := ensureCurrentPageRule(rule.ContentJSON); err != nil {
return rule, err
}
return rule, nil
}
if err := tx.First(&rule, record.RuleID).Error; err != nil {
if errors.Is(err, gorm.ErrRecordNotFound) {
return rule, serviceError(CodeRuleNotFound, "采集规则不存在或已删除")
}
return rule, internalError(err)
}
return rule, nil
}
func archiveCollectionAttempt(tx *gorm.DB, record models.CollectionTask, resetRequestID string, archivedAt time.Time) error {
counts := struct {
Dimensions int64 `json:"dimensions"`
ColorPrices int64 `json:"colorPrices"`
SKUs int64 `json:"skus"`
}{}
queries := []struct {
Model any
Count *int64
}{
{Model: &models.CollectionDimension{}, Count: &counts.Dimensions},
{Model: &models.CollectionColorPrice{}, Count: &counts.ColorPrices},
{Model: &models.CollectionSKU{}, Count: &counts.SKUs},
}
for _, query := range queries {
if err := tx.Model(query.Model).Where("task_id = ?", record.ID).Count(query.Count).Error; err != nil {
return internalError(err)
}
}
summary := struct {
Title *string `json:"title,omitempty"`
ShopName *string `json:"shopName,omitempty"`
MissingJSON *string `json:"missing,omitempty"`
DimensionCount int64 `json:"dimensionCount"`
ColorPriceCount int64 `json:"colorPriceCount"`
SKUCount int64 `json:"skuCount"`
}{
Title: record.Title, ShopName: record.ShopName, MissingJSON: record.MissingJSON,
DimensionCount: counts.Dimensions, ColorPriceCount: counts.ColorPrices, SKUCount: counts.SKUs,
}
raw, err := json.Marshal(summary)
if err != nil {
return internalError(err)
}
attemptNumber := record.AttemptNumber
if attemptNumber < 1 {
attemptNumber = 1
}
attempt := models.CollectionTaskAttempt{
TaskID: record.ID, AttemptNumber: attemptNumber, Source: record.Source, DeviceID: record.DeviceID,
RuleID: record.RuleID, RuleSnapshot: record.RuleSnapshot, Status: record.Status,
ErrorCode: record.ErrorCode, ErrorMessage: record.ErrorMessage, ResultSummaryJSON: string(raw),
StartedAt: record.StartedAt, FinishedAt: record.FinishedAt,
ArchivedByResetRequestID: resetRequestID, ArchivedAt: archivedAt,
}
if err := tx.Create(&attempt).Error; err != nil {
return internalError(err)
}
return nil
}
func ensureDeviceIdleForReset(tx *gorm.DB, deviceID, taskID uint64, now time.Time) error {
var busy int64
if err := tx.Model(&models.CollectionTask{}).
+11 -7
View File
@@ -71,12 +71,13 @@ type DetailSKU struct {
Complete bool `json:"complete"`
}
type DetailResponse struct {
Task AdminTaskItem `json:"task"`
Dimensions []DetailDimension `json:"dimensions"`
ColorPrices []models.CollectionColorPrice `json:"colorPrices"`
SKUs []DetailSKU `json:"skus"`
Missing []string `json:"missing"`
Replayed bool `json:"replayed,omitempty"`
Task AdminTaskItem `json:"task"`
Attempts []models.CollectionTaskAttempt `json:"attempts"`
Dimensions []DetailDimension `json:"dimensions"`
ColorPrices []models.CollectionColorPrice `json:"colorPrices"`
SKUs []DetailSKU `json:"skus"`
Missing []string `json:"missing"`
Replayed bool `json:"replayed,omitempty"`
}
func (service *Service) SubmitResult(ctx context.Context, taskID uint64, request ResultRequest, token string) (DetailResponse, error) {
@@ -397,7 +398,10 @@ func (service *Service) Detail(ctx context.Context, taskID uint64) (DetailRespon
if err := db.Where("task_id = ?", taskID).Order("sort_order,id").Find(&dimensions).Error; err != nil {
return DetailResponse{}, internalError(err)
}
result := DetailResponse{Task: item, Dimensions: make([]DetailDimension, 0), ColorPrices: make([]models.CollectionColorPrice, 0), SKUs: make([]DetailSKU, 0), Missing: make([]string, 0)}
result := DetailResponse{Task: item, Attempts: make([]models.CollectionTaskAttempt, 0), Dimensions: make([]DetailDimension, 0), ColorPrices: make([]models.CollectionColorPrice, 0), SKUs: make([]DetailSKU, 0), Missing: make([]string, 0)}
if err := db.Where("task_id = ?", taskID).Order("attempt_number,id").Find(&result.Attempts).Error; err != nil {
return DetailResponse{}, internalError(err)
}
valueLookup := map[uint64]struct{ Key, Value string }{}
for _, dimension := range dimensions {
var values []models.CollectionDimensionValue
+2 -1
View File
@@ -52,6 +52,7 @@ type ActionRequest struct {
type TaskPayload struct {
TaskID uint64 `json:"taskId"`
AttemptNumber int `json:"attemptNumber"`
PDDProductID *uint64 `json:"pddProductId"`
URLSnapshot string `json:"urlSnapshot"`
GoodsIDSnapshot string `json:"goodsIdSnapshot"`
@@ -290,7 +291,7 @@ func (service *Service) payload(record models.CollectionTask, replayed bool) (Ta
timeout = DefaultTaskTimeout
}
return TaskPayload{
TaskID: record.ID, PDDProductID: record.PDDProductID,
TaskID: record.ID, AttemptNumber: record.AttemptNumber, PDDProductID: record.PDDProductID,
URLSnapshot: record.URLSnapshot, GoodsIDSnapshot: record.GoodsIDSnapshot,
Source: record.Source, ReplacementOriginType: pointerValue(record.ReplacementOriginType),
RuleID: record.RuleID, RuleSnapshot: rule, TimeoutSeconds: timeout,
@@ -0,0 +1,29 @@
package version_local
import (
"runtime"
goautomigrations "go-admin/app/goauto/migrations"
"go-admin/cmd/migrate/migration"
common "go-admin/common/models"
"gorm.io/gorm"
)
// #155 adds collection attempt auditing and the current attempt counter.
func init() {
_, fileName, _, _ := runtime.Caller(0)
migration.Migrate.SetVersion(migration.GetFilename(fileName), migrateCollectionTaskAttempt)
}
func migrateCollectionTaskAttempt(db *gorm.DB, version string) error {
return db.Transaction(func(tx *gorm.DB) error {
if err := goautomigrations.Migrate(tx); err != nil {
return err
}
if err := tx.Model(&common.Migration{}).Create(&common.Migration{Version: version}).Error; err != nil {
return err
}
return nil
})
}
@@ -0,0 +1,71 @@
package version_local
import (
"testing"
goautomigrations "go-admin/app/goauto/migrations"
"go-admin/app/goauto/models"
common "go-admin/common/models"
"gorm.io/driver/sqlite"
"gorm.io/gorm"
)
type collectionTaskBefore155 struct {
ID uint64 `gorm:"primaryKey;autoIncrement"`
RuleID uint64 `gorm:"not null;index"`
Status string `gorm:"size:24;not null"`
URLSnapshot string `gorm:"type:text;not null"`
GoodsIDSnapshot string `gorm:"size:32;not null"`
RuleSnapshot string `gorm:"type:text;not null"`
}
func (collectionTaskBefore155) TableName() string { return "collection_task" }
func TestMigrateCollectionTaskAttemptUpgradesExistingSchema(t *testing.T) {
db, err := gorm.Open(sqlite.Open("file:collection-task-attempt-upgrade?mode=memory&cache=shared"), &gorm.Config{})
if err != nil {
t.Fatal(err)
}
if err = db.AutoMigrate(&common.Migration{}); err != nil {
t.Fatal(err)
}
if err = db.AutoMigrate(&collectionTaskBefore155{}); err != nil {
t.Fatal(err)
}
if err = db.Create(&collectionTaskBefore155{RuleID: 1, Status: "failed", URLSnapshot: "url", GoodsIDSnapshot: "12345", RuleSnapshot: "{}"}).Error; err != nil {
t.Fatal(err)
}
if db.Migrator().HasTable(&models.CollectionTaskAttempt{}) {
t.Fatal("upgrade precondition failed: collection attempt table already exists")
}
if db.Migrator().HasColumn(&models.CollectionTask{}, "AttemptNumber") {
t.Fatal("upgrade precondition failed: collection task attempt_number already exists")
}
const version = "1787983500000"
if err = migrateCollectionTaskAttempt(db, version); err != nil {
t.Fatal(err)
}
if !db.Migrator().HasTable(&models.CollectionTaskAttempt{}) {
t.Fatal("collection attempt table was not created")
}
if !db.Migrator().HasColumn(&models.CollectionTask{}, "AttemptNumber") {
t.Fatal("collection task attempt_number was not created")
}
var attemptNumber int
if err = db.Table("collection_task").Select("attempt_number").Where("id = 1").Scan(&attemptNumber).Error; err != nil || attemptNumber != 1 {
t.Fatalf("existing collection task attempt_number=%d err=%v", attemptNumber, err)
}
if !db.Migrator().HasIndex(&models.CollectionTaskAttempt{}, "ux_collection_attempt_number") ||
!db.Migrator().HasIndex(&models.CollectionTaskAttempt{}, "ux_collection_attempt_reset_request") {
t.Fatal("collection attempt unique indexes were not created")
}
var count int64
if err = db.Model(&common.Migration{}).Where("version = ?", version).Count(&count).Error; err != nil || count != 1 {
t.Fatalf("migration version count=%d err=%v", count, err)
}
if err = goautomigrations.VerifyTables(db); err != nil {
t.Fatalf("schema verification failed after upgrade: %v", err)
}
}
@@ -13,6 +13,7 @@
</el-form>
<el-table v-loading="loading" :data="tasks" border stripe empty-text="暂无采集任务">
<el-table-column label="任务 ID" prop="id" width="100" />
<el-table-column label="执行次数" width="90" align="center"><template #default="{ row }">第 {{ row.attemptNumber || 1 }} 次</template></el-table-column>
<el-table-column label="Goods ID" prop="goodsIdSnapshot" min-width="160" />
<el-table-column label="规则" prop="ruleName" min-width="150" show-overflow-tooltip />
<el-table-column label="设备" min-width="150"><template #default="{ row }">{{ row.deviceName || '空闲设备领取' }}</template></el-table-column>
@@ -40,7 +41,8 @@
<el-drawer v-model="detail.open" title="采集任务详情" size="720px">
<div v-if="detail.data" v-loading="detail.loading" class="detail-body">
<el-descriptions :column="2" border>
<el-descriptions-item label="任务 ID">{{ detail.data.task.id }}</el-descriptions-item><el-descriptions-item label="状态"><el-tag :type="statusType(detail.data.task.status)">{{ statusLabel(detail.data.task.status) }}</el-tag></el-descriptions-item>
<el-descriptions-item label="任务 ID">{{ detail.data.task.id }}</el-descriptions-item><el-descriptions-item label="当前执行">第 {{ detail.data.task.attemptNumber || 1 }} 次</el-descriptions-item>
<el-descriptions-item label="状态"><el-tag :type="statusType(detail.data.task.status)">{{ statusLabel(detail.data.task.status) }}</el-tag></el-descriptions-item>
<el-descriptions-item label="Goods ID">{{ detail.data.task.goodsIdSnapshot }}</el-descriptions-item><el-descriptions-item label="设备">{{ detail.data.task.deviceName || '未指定' }}</el-descriptions-item>
<el-descriptions-item label="标题" :span="2">{{ detail.data.task.title || '未采集' }}</el-descriptions-item><el-descriptions-item label="店铺">{{ detail.data.task.shopName || '未采集' }}</el-descriptions-item><el-descriptions-item label="销量">{{ detail.data.task.salesText || '未采集' }}</el-descriptions-item>
<el-descriptions-item label="评价数">{{ detail.data.task.reviewCount ?? '未采集' }}</el-descriptions-item><el-descriptions-item label="错误">{{ detail.data.task.errorCode ? `${detail.data.task.errorCode}:${detail.data.task.errorMessage}` : '无' }}</el-descriptions-item>
@@ -49,6 +51,17 @@
<h3>颜色价格</h3><el-table :data="detail.data.colorPrices" border empty-text="暂无颜色价格"><el-table-column prop="color" label="颜色" /><el-table-column label="价格" width="140"><template #default="{ row }">¥{{ (row.priceCent / 100).toFixed(2) }}</template></el-table-column></el-table>
<h3>SKU</h3><el-table :data="detail.data.skus" border empty-text="暂无 SKU"><el-table-column label="规格"><template #default="{ row }">{{ specsText(row.specs) }}</template></el-table-column><el-table-column label="价格" width="120"><template #default="{ row }">¥{{ (row.priceCent / 100).toFixed(2) }}</template></el-table-column><el-table-column label="完整" width="90"><template #default="{ row }">{{ row.complete ? '是' : '否' }}</template></el-table-column></el-table>
<el-alert v-if="detail.data.missing?.length" :title="`缺失项:${detail.data.missing.join('、')}`" type="warning" :closable="false" show-icon class="missing-alert" />
<template v-if="detail.data.attempts?.length">
<h3>历史执行记录</h3>
<el-table :data="detail.data.attempts" border>
<el-table-column type="expand"><template #default="{ row }"><pre class="snapshot-json" tabindex="0">{{ formatRuleSnapshot(row.ruleSnapshot) }}</pre></template></el-table-column>
<el-table-column label="次数" width="80"><template #default="{ row }">第 {{ row.attemptNumber }} 次</template></el-table-column>
<el-table-column label="状态" width="110"><template #default="{ row }"><el-tag :type="statusType(row.status)">{{ statusLabel(row.status) }}</el-tag></template></el-table-column>
<el-table-column prop="ruleId" label="规则 ID" width="90" />
<el-table-column label="错误"><template #default="{ row }">{{ row.errorCode ? `${row.errorCode}:${row.errorMessage || ''}` : '无' }}</template></el-table-column>
<el-table-column label="结束时间" width="180"><template #default="{ row }">{{ row.finishedAt ? parseTime(row.finishedAt) : '未记录' }}</template></el-table-column>
</el-table>
</template>
<el-collapse class="snapshot-collapse"><el-collapse-item title="查看本任务规则快照" name="rule-snapshot"><pre class="snapshot-json" tabindex="0">{{ formatRuleSnapshot(detail.data.task.ruleSnapshot) }}</pre></el-collapse-item></el-collapse>
</div>
</el-drawer>
@@ -77,7 +90,7 @@ export default {
async openCreate() { this.form = { pddProductId: null, ruleId: null, deviceId: null }; this.createDialog.open = true; const [p, r, d] = await Promise.all([listPddProducts({ page: 1, pageSize: 100 }), listCollectionRules({ page: 1, pageSize: 100 }), listDevices({ page: 1, pageSize: 100 })]); this.options = { products: p.data.items, rules: r.data.items, devices: d.data.items }; this.$nextTick(() => this.$refs.taskForm?.clearValidate()) },
async createTask() { if (!await this.$refs.taskForm.validate().catch(() => false)) return; this.createDialog.saving = true; try { await createCollectionTask({ requestId: window.crypto.randomUUID(), pddProductId: this.form.pddProductId, ruleId: this.form.ruleId, deviceId: this.form.deviceId || null }); ElMessage.success('任务已创建'); this.createDialog.open = false; await this.getList() } finally { this.createDialog.saving = false } },
async showDetail(row) { this.detail = { open: true, loading: true, data: null }; try { const r = await getCollectionTask(row.id); this.detail.data = r.data } finally { this.detail.loading = false } },
async confirmReset(row) { await ElMessageBox.confirm('重置会清除本次采集结果并重新进入待执行,是否继续?', '重置任务', { type: 'warning' }); await resetCollectionTask(row.id, { requestId: window.crypto.randomUUID() }); ElMessage.success('任务已重置'); await this.getList() },
async confirmReset(row) { await ElMessageBox.confirm('将复用原任务并按最新有效规则重新采集;当前结果会归档到历史执行记录,不会创建新任务。是否继续?', '重新采集', { type: 'warning' }); const response = await resetCollectionTask(row.id, { requestId: window.crypto.randomUUID() }); ElMessage.success(`任务已进入第 ${response.data.attemptNumber || 1} 次采集`); await this.getList() },
async confirmDelete(row) { await ElMessageBox.confirm('仅删除这条失败任务,是否继续?', '删除失败任务', { type: 'warning' }); await deleteCollectionTask(row.id, { requestId: window.crypto.randomUUID() }); ElMessage.success('失败任务已删除'); await this.getList() }
}
}