From 2372ab22806b4edcf2e1386cb2daaf3c01655787 Mon Sep 17 00:00:00 2001 From: QiuSW <105186638@qq.com> Date: Tue, 28 Jul 2026 13:23:18 +0800 Subject: [PATCH] feat(t216): deliver recoverable device order commands --- android-buyer/app/build.gradle.kts | 4 +- .../java/com/roubao/autopilot/MainActivity.kt | 68 ++- .../procurement/OrderCommandIntegrity.kt | 133 ++++++ .../OrderCommandSynchronization.kt | 43 ++ .../procurement/ProcurementApiClient.kt | 109 +++++ .../procurement/ProcurementModels.kt | 41 +- .../procurement/ProcurementRepository.kt | 113 ++++- .../procurement/ProcurementSecureStore.kt | 81 +++- .../autopilot/ui/screens/ProcurementScreen.kt | 30 +- .../autopilot/ui/screens/SearchProbeScreen.kt | 12 + .../vlm/CandidateEvaluationModels.kt | 1 + .../procurement/OrderCommandIntegrityTest.kt | 37 ++ .../OrderCommandSynchronizationTest.kt | 134 ++++++ .../procurement/ProcurementApiClientTest.kt | 103 +++++ backend-api/cmd/api/main.go | 9 + backend-api/internal/domain/task.go | 26 ++ .../migration/claims_migration_test.go | 27 +- .../platform/migration/runner_test.go | 9 +- .../repository/sqlite/auth_repository_test.go | 9 +- .../sqlite/device_order_command_repository.go | 404 ++++++++++++++++++ .../sqlite/order_authorization_repository.go | 14 +- .../transport/httpapi/admin_handlers.go | 3 + .../transport/httpapi/admin_handlers_test.go | 3 + .../transport/httpapi/device_handlers.go | 117 ++++- .../transport/httpapi/device_handlers_test.go | 333 +++++++++++++++ .../usecase/device_order_command_service.go | 336 +++++++++++++++ .../device_order_command_service_test.go | 47 ++ .../00009_device_order_commands.sql | 175 ++++++++ docs/00-ai-start-here.md | 3 +- docs/08-interaction-checklist.md | 2 +- docs/current-state.md | 28 +- docs/tasks/T-216.md | 26 +- 32 files changed, 2415 insertions(+), 65 deletions(-) create mode 100644 android-buyer/app/src/main/java/com/roubao/autopilot/procurement/OrderCommandIntegrity.kt create mode 100644 android-buyer/app/src/main/java/com/roubao/autopilot/procurement/OrderCommandSynchronization.kt create mode 100644 android-buyer/app/src/test/java/com/roubao/autopilot/procurement/OrderCommandIntegrityTest.kt create mode 100644 android-buyer/app/src/test/java/com/roubao/autopilot/procurement/OrderCommandSynchronizationTest.kt create mode 100644 backend-api/internal/repository/sqlite/device_order_command_repository.go create mode 100644 backend-api/internal/usecase/device_order_command_service.go create mode 100644 backend-api/internal/usecase/device_order_command_service_test.go create mode 100644 backend-api/migrations/00009_device_order_commands.sql diff --git a/android-buyer/app/build.gradle.kts b/android-buyer/app/build.gradle.kts index 0fa962e..2b26879 100644 --- a/android-buyer/app/build.gradle.kts +++ b/android-buyer/app/build.gradle.kts @@ -11,8 +11,8 @@ android { applicationId = "com.roubao.autopilot" minSdk = 26 targetSdk = 34 - versionCode = 17 - versionName = "1.4.12" + versionCode = 18 + versionName = "1.4.13" vectorDrawables { useSupportLibrary = true diff --git a/android-buyer/app/src/main/java/com/roubao/autopilot/MainActivity.kt b/android-buyer/app/src/main/java/com/roubao/autopilot/MainActivity.kt index ed95b93..dbaea63 100644 --- a/android-buyer/app/src/main/java/com/roubao/autopilot/MainActivity.kt +++ b/android-buyer/app/src/main/java/com/roubao/autopilot/MainActivity.kt @@ -106,6 +106,7 @@ import com.roubao.autopilot.procurement.ExecutionEvidenceKind import com.roubao.autopilot.procurement.ExecutionMode import com.roubao.autopilot.procurement.ExecutionProvenanceSnapshot import com.roubao.autopilot.procurement.ProcurementRepository +import com.roubao.autopilot.procurement.ProcurementPhase import com.roubao.autopilot.procurement.RankedExecutionCandidateDraft import com.roubao.autopilot.procurement.CandidateHumanReviewDraft import com.roubao.autopilot.task.RequirementProbeFixture @@ -296,6 +297,16 @@ class MainActivity : ComponentActivity() { LaunchedEffect(procurementState.task?.id) { resetTaskBoundProbeState() } + LaunchedEffect(procurementState.phase) { + if ( + procurementState.phase == + ProcurementPhase.WAITING_ADMIN_CONFIRMATION || + procurementState.phase == ProcurementPhase.ORDER_AUTHORIZED + ) { + candidateEvaluationState.value = + CandidateEvaluationState.AWAITING_ADMIN_CONFIRMATION + } + } // 监听跳转事件 LaunchedEffect(navigateToRecord, recordId) { @@ -436,7 +447,16 @@ class MainActivity : ComponentActivity() { .map { it.candidate.ordinal }, hasBudget = procurementState.task?.maxBudget != null, + candidateCollectionAllowed = + procurementState.task == null || + procurementState.phase == + ProcurementPhase.RUNNING, canStartCandidateEvaluation = + ( + procurementState.task == null || + procurementState.phase == + ProcurementPhase.RUNNING + ) && extractionState == RequirementProbeState.READY && extraction != null && evidenceRequirement == extraction && @@ -598,6 +618,10 @@ class MainActivity : ComponentActivity() { Toast.makeText(this, "设备检查存在阻塞项", Toast.LENGTH_SHORT).show() return } + if (!procurementRepository.canCollectCandidates) { + Toast.makeText(this, "当前任务已进入后台确认阶段", Toast.LENGTH_SHORT).show() + return + } if (searchProbeJob?.isActive == true) { return } @@ -773,6 +797,10 @@ class MainActivity : ComponentActivity() { } private fun startRequirementProbe() { + if (!procurementRepository.canCollectCandidates) { + Toast.makeText(this, "当前任务已进入后台确认阶段", Toast.LENGTH_SHORT).show() + return + } if (requirementProbeJob?.isActive == true) { return } @@ -900,6 +928,10 @@ class MainActivity : ComponentActivity() { } private fun startCandidateEvaluation() { + if (!procurementRepository.canCollectCandidates) { + Toast.makeText(this, "当前任务已进入后台确认阶段", Toast.LENGTH_SHORT).show() + return + } if (candidateEvaluationJob?.isActive == true) { return } @@ -1187,16 +1219,24 @@ class MainActivity : ComponentActivity() { } } candidateReviewBatch.value = result.batch - candidateEvaluationState.value = when ( - result.batch.conclusion - ) { - CandidateBatchConclusion.SUGGESTED -> - CandidateEvaluationState.AWAITING_CONFIRMATION - CandidateBatchConclusion.NO_MATCH -> - CandidateEvaluationState.NO_MATCH - CandidateBatchConclusion.MANUAL_REQUIRED -> - CandidateEvaluationState.MANUAL_REVIEW - } + candidateEvaluationState.value = + if ( + procurementRepository.currentProbeTask() != null && + taskCandidateDrafts.value.isNotEmpty() + ) { + CandidateEvaluationState + .AWAITING_ADMIN_CONFIRMATION + } else { + when (result.batch.conclusion) { + CandidateBatchConclusion.SUGGESTED -> + CandidateEvaluationState + .AWAITING_CONFIRMATION + CandidateBatchConclusion.NO_MATCH -> + CandidateEvaluationState.NO_MATCH + CandidateBatchConclusion.MANUAL_REQUIRED -> + CandidateEvaluationState.MANUAL_REVIEW + } + } } is CandidateEvaluationResult.Failed -> { setCandidateEvaluationFailure(result.code) @@ -1281,7 +1321,10 @@ class MainActivity : ComponentActivity() { return } if (taskCandidateDrafts.value.isNotEmpty()) { - candidateEvaluationState.value = CandidateEvaluationState.MANUAL_REVIEW + candidateEvaluationState.value = + CandidateEvaluationState.AWAITING_ADMIN_CONFIRMATION + } else { + candidateEvaluationState.value = CandidateEvaluationState.NO_MATCH } } @@ -1337,7 +1380,8 @@ class MainActivity : ComponentActivity() { if ( candidateEvaluationState.value == CandidateEvaluationState.HUMAN_REJECTED && - procurementRepository.currentProbeTask() != null + procurementRepository.currentProbeTask() != null && + taskCandidateDrafts.value.isEmpty() ) { lifecycleScope.launch { procurementRepository.completeExecution(review = review) diff --git a/android-buyer/app/src/main/java/com/roubao/autopilot/procurement/OrderCommandIntegrity.kt b/android-buyer/app/src/main/java/com/roubao/autopilot/procurement/OrderCommandIntegrity.kt new file mode 100644 index 0000000..84c632d --- /dev/null +++ b/android-buyer/app/src/main/java/com/roubao/autopilot/procurement/OrderCommandIntegrity.kt @@ -0,0 +1,133 @@ +package com.roubao.autopilot.procurement + +import java.security.MessageDigest + +object OrderCommandIntegrity { + const val SCHEMA_VERSION = 1 + const val TYPE = "CREATE_PENDING_ORDER" + private val sha256Pattern = Regex("[0-9a-f]{64}") + + fun validate( + command: PendingOrderCommand, + task: RemotePurchaseTask, + execution: RunningExecution + ) { + require(command.schemaVersion == SCHEMA_VERSION) { + "后台订单命令版本不受支持" + } + require(command.type == TYPE) { "后台订单命令类型不受支持" } + require( + command.id.isNotBlank() && + command.authorizationVersion >= 1 && + command.taskId == task.id && + command.executionId == execution.id && + command.taskContentSha256 == ExecutionTaskHash.sha256(task) && + command.originalSku == task.sku && + command.quantity == task.quantity + ) { "后台订单命令与当前任务不一致" } + require(command.candidate.observedOrdinal in 1..5) { + "后台订单命令候选编号无效" + } + require(command.candidate.title.isNotBlank()) { + "后台订单命令候选标题缺失" + } + require( + listOf( + command.candidate.candidateKey, + command.candidate.cardSignature, + command.candidate.detailSignature, + command.candidate.detailEvidenceSha256, + command.candidate.specificationEvidenceSha256, + command.commandSha256 + ).all(sha256Pattern::matches) + ) { "后台订单命令指纹无效" } + require( + command.authorizationStatus == "DELIVERED" || + command.authorizationStatus == "ACKNOWLEDGED" + ) { "后台订单命令状态无效" } + require(sha256(command) == command.commandSha256) { + "后台订单命令完整性校验失败" + } + } + + fun sha256(command: PendingOrderCommand): String { + val payload = buildString { + append('{') + property("id", command.id, first = true) + property("schema_version", command.schemaVersion) + property("type", command.type) + property("authorization_version", command.authorizationVersion) + property("task_id", command.taskId) + property("execution_id", command.executionId) + property("task_content_sha256", command.taskContentSha256) + property("original_sku", command.originalSku) + property("quantity", command.quantity) + append(",\"candidate\":{") + property("candidate_key", command.candidate.candidateKey, first = true) + property("observed_ordinal", command.candidate.observedOrdinal) + property("title", command.candidate.title) + property("sku_text", command.candidate.skuText) + property("price_text", command.candidate.priceText) + property("card_signature", command.candidate.cardSignature) + property("detail_signature", command.candidate.detailSignature) + property( + "detail_evidence_sha256", + command.candidate.detailEvidenceSha256 + ) + property( + "specification_evidence_sha256", + command.candidate.specificationEvidenceSha256 + ) + append("}}") + } + return MessageDigest.getInstance("SHA-256") + .digest(payload.toByteArray(Charsets.UTF_8)) + .joinToString("") { "%02x".format(it) } + } + + private fun StringBuilder.property( + name: String, + value: String, + first: Boolean = false + ) { + if (!first) append(',') + append('"').append(name).append("\":") + appendJsonString(value) + } + + private fun StringBuilder.property( + name: String, + value: Int, + first: Boolean = false + ) { + if (!first) append(',') + append('"').append(name).append("\":").append(value) + } + + private fun StringBuilder.appendJsonString(value: String) { + append('"') + value.forEach { character -> + when (character) { + '"' -> append("\\\"") + '\\' -> append("\\\\") + '\b' -> append("\\b") + '\u000C' -> append("\\f") + '\n' -> append("\\n") + '\r' -> append("\\r") + '\t' -> append("\\t") + '<' -> append("\\u003c") + '>' -> append("\\u003e") + '&' -> append("\\u0026") + '\u2028' -> append("\\u2028") + '\u2029' -> append("\\u2029") + else -> if (character.code < 0x20) { + append("\\u") + append(character.code.toString(16).padStart(4, '0')) + } else { + append(character) + } + } + } + append('"') + } +} diff --git a/android-buyer/app/src/main/java/com/roubao/autopilot/procurement/OrderCommandSynchronization.kt b/android-buyer/app/src/main/java/com/roubao/autopilot/procurement/OrderCommandSynchronization.kt new file mode 100644 index 0000000..70011cf --- /dev/null +++ b/android-buyer/app/src/main/java/com/roubao/autopilot/procurement/OrderCommandSynchronization.kt @@ -0,0 +1,43 @@ +package com.roubao.autopilot.procurement + +object OrderCommandSynchronization { + suspend fun synchronize( + existing: PendingOrderCommand?, + task: RemotePurchaseTask, + execution: RunningExecution, + pull: suspend () -> PendingOrderCommand?, + persist: (PendingOrderCommand) -> Unit, + acknowledge: suspend ( + PendingOrderCommand, + String + ) -> OrderCommandAcknowledgement, + newIdempotencyKey: () -> String + ): PendingOrderCommand? { + var command = existing + if (command == null) { + val delivered = pull() ?: return null + OrderCommandIntegrity.validate(delivered, task, execution) + command = delivered.copy( + acknowledgementIdempotencyKey = newIdempotencyKey() + ) + persist(command) + } else { + OrderCommandIntegrity.validate(command, task, execution) + } + if (command.acknowledged) { + return command + } + val idempotencyKey = requireNotNull( + command.acknowledgementIdempotencyKey + ) { "后台订单命令确认键缺失" } + val acknowledgement = acknowledge(command, idempotencyKey) + require( + acknowledgement.commandId == command.id && + acknowledgement.status == "ACKNOWLEDGED" + ) { "后台订单命令确认响应无效" } + return command.copy( + authorizationStatus = acknowledgement.status, + acknowledged = true + ).also(persist) + } +} diff --git a/android-buyer/app/src/main/java/com/roubao/autopilot/procurement/ProcurementApiClient.kt b/android-buyer/app/src/main/java/com/roubao/autopilot/procurement/ProcurementApiClient.kt index 8157acb..ff54e76 100644 --- a/android-buyer/app/src/main/java/com/roubao/autopilot/procurement/ProcurementApiClient.kt +++ b/android-buyer/app/src/main/java/com/roubao/autopilot/procurement/ProcurementApiClient.kt @@ -84,6 +84,22 @@ interface ProcurementRemoteApi { idempotencyKey: String ): RemotePurchaseTask + suspend fun pullOrderCommand( + session: ProcurementSession, + task: RemotePurchaseTask, + execution: RunningExecution, + claimToken: String + ): PendingOrderCommand? + + suspend fun acknowledgeOrderCommand( + session: ProcurementSession, + task: RemotePurchaseTask, + execution: RunningExecution, + claimToken: String, + command: PendingOrderCommand, + idempotencyKey: String + ): OrderCommandAcknowledgement + suspend fun uploadExecutionOutboxItem( session: ProcurementSession, task: RemotePurchaseTask, @@ -395,6 +411,67 @@ class ProcurementApiClient( return parseTask(json.getJSONObject("task")) } + override suspend fun pullOrderCommand( + session: ProcurementSession, + task: RemotePurchaseTask, + execution: RunningExecution, + claimToken: String + ): PendingOrderCommand? = withContext(Dispatchers.IO) { + val payload = JSONObject() + .put("device_id", session.deviceId) + .put("execution_id", execution.id) + .put("claim_generation", task.claimGeneration) + client.newCall( + authorizedRequest( + session, + "/api/v1/tasks/${task.id}/commands/next" + ) + .header(CLAIM_TOKEN_HEADER, claimToken) + .post(payload.jsonBody()) + .build() + ).execute().use { response -> + if (response.code == 204) { + return@withContext null + } + parseOrderCommand( + responseJsonOrThrow( + response.code, + response.body.readBoundedJson(response.code) + ) + ) + } + } + + override suspend fun acknowledgeOrderCommand( + session: ProcurementSession, + task: RemotePurchaseTask, + execution: RunningExecution, + claimToken: String, + command: PendingOrderCommand, + idempotencyKey: String + ): OrderCommandAcknowledgement { + val payload = JSONObject() + .put("device_id", session.deviceId) + .put("execution_id", execution.id) + .put("claim_generation", task.claimGeneration) + .put("command_sha256", command.commandSha256) + val json = executeJson( + authorizedRequest( + session, + "/api/v1/tasks/${task.id}/commands/${command.id}/ack" + ) + .header(CLAIM_TOKEN_HEADER, claimToken) + .header(IDEMPOTENCY_HEADER, idempotencyKey) + .post(payload.jsonBody()) + .build() + ) + return OrderCommandAcknowledgement( + commandId = json.getString("command_id"), + status = json.getString("status"), + replayed = json.getBoolean("replayed") + ) + } + private suspend fun executeTransition( session: ProcurementSession, task: RemotePurchaseTask, @@ -466,6 +543,8 @@ class ProcurementApiClient( "TASK_CLAIM_INVALID" -> "任务归属校验失败" "TASK_CLAIM_EXPIRED" -> "任务授权已过期" "TASK_VERSION_CONFLICT" -> "任务已被后台更新,请先同步" + "ORDER_COMMAND_CONFLICT" -> "后台订单命令已变化,自动化保持停止" + "ORDER_COMMAND_NOT_DELIVERABLE" -> "后台订单命令当前不可领取" else -> if (statusCode >= 500) "后台暂时不可用" else "后台拒绝了请求" } @@ -521,6 +600,36 @@ class ProcurementApiClient( currency = json.optString("currency", "CNY") ) + private fun parseOrderCommand(json: JSONObject): PendingOrderCommand { + val candidate = json.getJSONObject("candidate") + return PendingOrderCommand( + id = json.getString("id"), + schemaVersion = json.getInt("schema_version"), + type = json.getString("type"), + authorizationVersion = json.getInt("authorization_version"), + taskId = json.getString("task_id"), + executionId = json.getString("execution_id"), + taskContentSha256 = json.getString("task_content_sha256"), + originalSku = json.getString("original_sku"), + quantity = json.getInt("quantity"), + candidate = OrderCommandCandidate( + candidateKey = candidate.getString("candidate_key"), + observedOrdinal = candidate.getInt("observed_ordinal"), + title = candidate.getString("title"), + skuText = candidate.getString("sku_text"), + priceText = candidate.getString("price_text"), + cardSignature = candidate.getString("card_signature"), + detailSignature = candidate.getString("detail_signature"), + detailEvidenceSha256 = + candidate.getString("detail_evidence_sha256"), + specificationEvidenceSha256 = + candidate.getString("specification_evidence_sha256") + ), + commandSha256 = json.getString("command_sha256"), + authorizationStatus = json.getString("authorization_status") + ) + } + private fun JSONObject.jsonBody() = toString().toRequestBody(JSON_MEDIA_TYPE) diff --git a/android-buyer/app/src/main/java/com/roubao/autopilot/procurement/ProcurementModels.kt b/android-buyer/app/src/main/java/com/roubao/autopilot/procurement/ProcurementModels.kt index 15a3037..7b44b79 100644 --- a/android-buyer/app/src/main/java/com/roubao/autopilot/procurement/ProcurementModels.kt +++ b/android-buyer/app/src/main/java/com/roubao/autopilot/procurement/ProcurementModels.kt @@ -9,6 +9,8 @@ enum class ProcurementPhase { IDLE, CLAIMED, RUNNING, + WAITING_ADMIN_CONFIRMATION, + ORDER_AUTHORIZED, AUTHORIZATION_EXPIRED } @@ -84,7 +86,43 @@ data class PersistedProcurementState( val deviceToken: String? = null, val claim: ClaimContext? = null, val execution: RunningExecution? = null, - val outbox: List = emptyList() + val outbox: List = emptyList(), + val orderCommand: PendingOrderCommand? = null +) + +data class OrderCommandCandidate( + val candidateKey: String, + val observedOrdinal: Int, + val title: String, + val skuText: String, + val priceText: String, + val cardSignature: String, + val detailSignature: String, + val detailEvidenceSha256: String, + val specificationEvidenceSha256: String +) + +data class PendingOrderCommand( + val id: String, + val schemaVersion: Int, + val type: String, + val authorizationVersion: Int, + val taskId: String, + val executionId: String, + val taskContentSha256: String, + val originalSku: String, + val quantity: Int, + val candidate: OrderCommandCandidate, + val commandSha256: String, + val authorizationStatus: String, + val acknowledgementIdempotencyKey: String? = null, + val acknowledged: Boolean = false +) + +data class OrderCommandAcknowledgement( + val commandId: String, + val status: String, + val replayed: Boolean ) enum class ExecutionMode { @@ -125,6 +163,7 @@ data class ProcurementUiState( val task: RemotePurchaseTask? = null, val referenceImagePath: String? = null, val execution: RunningExecution? = null, + val orderCommand: PendingOrderCommand? = null, val authenticationRequired: Boolean = false, val busy: Boolean = false, val backendOnline: Boolean? = null, diff --git a/android-buyer/app/src/main/java/com/roubao/autopilot/procurement/ProcurementRepository.kt b/android-buyer/app/src/main/java/com/roubao/autopilot/procurement/ProcurementRepository.kt index f987e91..e17e383 100644 --- a/android-buyer/app/src/main/java/com/roubao/autopilot/procurement/ProcurementRepository.kt +++ b/android-buyer/app/src/main/java/com/roubao/autopilot/procurement/ProcurementRepository.kt @@ -44,6 +44,19 @@ class ProcurementRepository( get() = storageFailure == null && persisted.execution?.safetyStopped == false + val canCollectCandidates: Boolean + get() { + if (persisted.claim?.task == null) { + return true + } + val execution = persisted.execution ?: return false + return !execution.safetyStopped && + !execution.isExpired() && + execution.currentStep != WAITING_ADMIN_CONFIRMATION_STEP && + execution.currentStep != ORDER_AUTHORIZED_STEP && + persisted.orderCommand == null + } + suspend fun login(input: LoginInput): Boolean = operation { require(input.password.isNotEmpty()) { "请输入采购员密码" } val normalizedUrl = BackendEndpointPolicy.normalize( @@ -342,12 +355,29 @@ class ProcurementRepository( acknowledgeCancellation(session) ExecutionSyncDecision.STOP } else { + val orderCommand = if ( + result.task.status == "WAITING_CONFIRMATION" || + persisted.orderCommand != null + ) { + synchronizeOrderCommand( + session = session, + task = result.task, + execution = updatedExecution, + claimToken = claim.token + ) + } else { + null + } publish( backendOnline = true, - message = if (updatedExecution.safetyStopped) { - "后台状态已同步,任务仍保持安全停止" - } else { - "进度已同步" + message = when { + updatedExecution.safetyStopped -> + "后台状态已同步,任务仍保持安全停止" + orderCommand?.acknowledged == true -> + "后台确认商品已安全同步,等待执行下单" + result.task.status == "WAITING_CONFIRMATION" -> + "候选已回传,等待后台采购员确认" + else -> "进度已同步" } ) if (updatedExecution.safetyStopped) { @@ -356,10 +386,18 @@ class ProcurementRepository( ExecutionSyncDecision.CONTINUE } } - } catch (_: Exception) { + } catch (error: Exception) { + val localFailure = + error is IllegalArgumentException || + error is IllegalStateException publish( - backendOnline = false, - message = "后台暂时离线,将在授权截止前重试" + backendOnline = if (localFailure) true else false, + message = if (localFailure) { + null + } else { + "后台暂时离线,将在授权截止前重试" + }, + error = if (localFailure) userMessage(error) else null ) if (execution.safetyStopped || execution.isExpired()) { ExecutionSyncDecision.STOP @@ -477,6 +515,13 @@ class ProcurementRepository( payload = payload.toString() ) ) + if (candidates.isNotEmpty()) { + persisted = persisted.copy( + execution = execution.copy( + currentStep = WAITING_ADMIN_CONFIRMATION_STEP + ) + ) + } saveAndPublish( backendOnline = _uiState.value.backendOnline, message = "候选和截图已加入加密回传队列" @@ -635,6 +680,49 @@ class ProcurementRepository( ) } + private suspend fun synchronizeOrderCommand( + session: ProcurementSession, + task: RemotePurchaseTask, + execution: RunningExecution, + claimToken: String + ): PendingOrderCommand? { + val command = OrderCommandSynchronization.synchronize( + existing = persisted.orderCommand, + task = task, + execution = execution, + pull = { + api.pullOrderCommand( + session = session, + task = task, + execution = execution, + claimToken = claimToken + ) + }, + persist = { storedCommand -> + persisted = persisted.copy(orderCommand = storedCommand) + store.save(persisted) + }, + acknowledge = { storedCommand, idempotencyKey -> + api.acknowledgeOrderCommand( + session = session, + task = task, + execution = execution, + claimToken = claimToken, + command = storedCommand, + idempotencyKey = idempotencyKey + ) + }, + newIdempotencyKey = ::newOpaqueSecret + ) ?: return null + if (command.acknowledged) { + persisted = persisted.copy( + execution = execution.copy(currentStep = ORDER_AUTHORIZED_STEP) + ) + store.save(persisted) + } + return command + } + private suspend fun downloadAndPersistReference( session: ProcurementSession, claim: ClaimContext, @@ -687,7 +775,8 @@ class ProcurementRepository( persisted = persisted.copy( claim = null, execution = null, - outbox = emptyList() + outbox = emptyList(), + orderCommand = null ) } @@ -1010,6 +1099,10 @@ class ProcurementRepository( execution != null && (execution.safetyStopped || execution.isExpired()) -> ProcurementPhase.AUTHORIZATION_EXPIRED + persisted.orderCommand?.acknowledged == true -> + ProcurementPhase.ORDER_AUTHORIZED + execution?.currentStep == WAITING_ADMIN_CONFIRMATION_STEP -> + ProcurementPhase.WAITING_ADMIN_CONFIRMATION execution != null -> ProcurementPhase.RUNNING claim?.task != null -> ProcurementPhase.CLAIMED session?.isValid() == true -> ProcurementPhase.IDLE @@ -1031,6 +1124,7 @@ class ProcurementRepository( task = claim?.task, referenceImagePath = imagePath, execution = execution, + orderCommand = persisted.orderCommand, authenticationRequired = session?.isValid() != true ) } @@ -1068,6 +1162,9 @@ class ProcurementRepository( private const val REFERENCE_DIRECTORY = "procurement" private const val OUTBOX_DIRECTORY = "procurement-outbox" private const val CONTROLLED_WORKFLOW_STEP = "CONTROLLED_WORKFLOW" + private const val WAITING_ADMIN_CONFIRMATION_STEP = + "WAITING_ADMIN_CONFIRMATION" + private const val ORDER_AUTHORIZED_STEP = "ORDER_AUTHORIZED" private const val SAFE_STOPPED_STEP = "SAFE_STOPPED" private const val MAX_REFERENCE_DIMENSION = 4_096 private const val MAX_REFERENCE_PIXELS = 20_000_000L diff --git a/android-buyer/app/src/main/java/com/roubao/autopilot/procurement/ProcurementSecureStore.kt b/android-buyer/app/src/main/java/com/roubao/autopilot/procurement/ProcurementSecureStore.kt index e3b61f5..7678b7b 100644 --- a/android-buyer/app/src/main/java/com/roubao/autopilot/procurement/ProcurementSecureStore.kt +++ b/android-buyer/app/src/main/java/com/roubao/autopilot/procurement/ProcurementSecureStore.kt @@ -134,6 +134,9 @@ class ProcurementSecureStore(context: Context) : ProcurementStateStore { } } ) + state.orderCommand?.let { command -> + put("order_command", encodeOrderCommand(command)) + } } private fun decodeState(json: JSONObject): PersistedProcurementState = @@ -215,9 +218,85 @@ class ProcurementSecureStore(context: Context) : ProcurementStateStore { ) } } - } ?: emptyList() + } ?: emptyList(), + orderCommand = json.optionalObject("order_command")?.let { + decodeOrderCommand(it) + } ) + private fun encodeOrderCommand(command: PendingOrderCommand): JSONObject = + JSONObject().apply { + put("id", command.id) + put("schema_version", command.schemaVersion) + put("type", command.type) + put("authorization_version", command.authorizationVersion) + put("task_id", command.taskId) + put("execution_id", command.executionId) + put("task_content_sha256", command.taskContentSha256) + put("original_sku", command.originalSku) + put("quantity", command.quantity) + put( + "candidate", + JSONObject().apply { + put("candidate_key", command.candidate.candidateKey) + put("observed_ordinal", command.candidate.observedOrdinal) + put("title", command.candidate.title) + put("sku_text", command.candidate.skuText) + put("price_text", command.candidate.priceText) + put("card_signature", command.candidate.cardSignature) + put("detail_signature", command.candidate.detailSignature) + put( + "detail_evidence_sha256", + command.candidate.detailEvidenceSha256 + ) + put( + "specification_evidence_sha256", + command.candidate.specificationEvidenceSha256 + ) + } + ) + put("command_sha256", command.commandSha256) + put("authorization_status", command.authorizationStatus) + putNullable( + "acknowledgement_idempotency_key", + command.acknowledgementIdempotencyKey + ) + put("acknowledged", command.acknowledged) + } + + private fun decodeOrderCommand(json: JSONObject): PendingOrderCommand { + val candidate = json.getJSONObject("candidate") + return PendingOrderCommand( + id = json.getString("id"), + schemaVersion = json.getInt("schema_version"), + type = json.getString("type"), + authorizationVersion = json.getInt("authorization_version"), + taskId = json.getString("task_id"), + executionId = json.getString("execution_id"), + taskContentSha256 = json.getString("task_content_sha256"), + originalSku = json.getString("original_sku"), + quantity = json.getInt("quantity"), + candidate = OrderCommandCandidate( + candidateKey = candidate.getString("candidate_key"), + observedOrdinal = candidate.getInt("observed_ordinal"), + title = candidate.getString("title"), + skuText = candidate.getString("sku_text"), + priceText = candidate.getString("price_text"), + cardSignature = candidate.getString("card_signature"), + detailSignature = candidate.getString("detail_signature"), + detailEvidenceSha256 = + candidate.getString("detail_evidence_sha256"), + specificationEvidenceSha256 = + candidate.getString("specification_evidence_sha256") + ), + commandSha256 = json.getString("command_sha256"), + authorizationStatus = json.getString("authorization_status"), + acknowledgementIdempotencyKey = + json.optionalString("acknowledgement_idempotency_key"), + acknowledged = json.optBoolean("acknowledged", false) + ) + } + private fun encodeTask(task: RemotePurchaseTask): JSONObject = JSONObject().apply { put("id", task.id) diff --git a/android-buyer/app/src/main/java/com/roubao/autopilot/ui/screens/ProcurementScreen.kt b/android-buyer/app/src/main/java/com/roubao/autopilot/ui/screens/ProcurementScreen.kt index ab2f4f5..bc34b97 100644 --- a/android-buyer/app/src/main/java/com/roubao/autopilot/ui/screens/ProcurementScreen.kt +++ b/android-buyer/app/src/main/java/com/roubao/autopilot/ui/screens/ProcurementScreen.kt @@ -196,6 +196,8 @@ fun ProcurementScreen( } } ProcurementPhase.RUNNING, + ProcurementPhase.WAITING_ADMIN_CONFIRMATION, + ProcurementPhase.ORDER_AUTHORIZED, ProcurementPhase.AUTHORIZATION_EXPIRED -> { if (state.authenticationRequired) { item(key = "procurement-login") { @@ -452,6 +454,12 @@ private fun ExecutionDetails(state: ProcurementUiState) { Column(verticalArrangement = Arrangement.spacedBy(10.dp)) { Divider(color = colors.surfaceVariant) DetailRow("当前步骤", execution.currentStep) + state.orderCommand?.let { command -> + DetailRow("后台已选商品", command.candidate.title) + DetailRow("目标规格", command.candidate.skuText) + DetailRow("目标数量", command.quantity.toString()) + DetailRow("候选指纹", command.candidate.candidateKey) + } DetailRow( "后台连接", when (state.backendOnline) { @@ -472,8 +480,20 @@ private fun ExecutionDetails(state: ProcurementUiState) { } ) Text( - "订单提交保持禁用", - color = colors.warning, + when (state.phase) { + ProcurementPhase.WAITING_ADMIN_CONFIRMATION -> + "等待后台采购员确认商品" + ProcurementPhase.ORDER_AUTHORIZED -> + "商品授权已安全保存;本版本尚不执行下单" + else -> "订单提交保持禁用" + }, + color = if ( + state.phase == ProcurementPhase.ORDER_AUTHORIZED + ) { + colors.success + } else { + colors.warning + }, fontSize = 14.sp ) } @@ -506,7 +526,9 @@ private fun DetailRow(label: String, value: String) { private fun phaseColor(state: ProcurementUiState) = when (state.phase) { ProcurementPhase.AUTHORIZATION_EXPIRED -> BaoziTheme.colors.error - ProcurementPhase.RUNNING -> BaoziTheme.colors.success + ProcurementPhase.RUNNING, + ProcurementPhase.WAITING_ADMIN_CONFIRMATION, + ProcurementPhase.ORDER_AUTHORIZED -> BaoziTheme.colors.success else -> BaoziTheme.colors.textSecondary } @@ -516,5 +538,7 @@ private fun phaseLabel(state: ProcurementUiState): String = ProcurementPhase.IDLE -> "等待领取" ProcurementPhase.CLAIMED -> "已领取,等待开始" ProcurementPhase.RUNNING -> "受控流程运行中" + ProcurementPhase.WAITING_ADMIN_CONFIRMATION -> "等待后台确认商品" + ProcurementPhase.ORDER_AUTHORIZED -> "后台商品授权已同步" ProcurementPhase.AUTHORIZATION_EXPIRED -> "授权到期,已停止" } diff --git a/android-buyer/app/src/main/java/com/roubao/autopilot/ui/screens/SearchProbeScreen.kt b/android-buyer/app/src/main/java/com/roubao/autopilot/ui/screens/SearchProbeScreen.kt index c449131..42bcc21 100644 --- a/android-buyer/app/src/main/java/com/roubao/autopilot/ui/screens/SearchProbeScreen.kt +++ b/android-buyer/app/src/main/java/com/roubao/autopilot/ui/screens/SearchProbeScreen.kt @@ -111,6 +111,7 @@ fun SearchProbeScreen( candidateEvaluationFailureCode: CandidateEvaluationFailureCode?, candidateOrdinals: List, hasBudget: Boolean, + candidateCollectionAllowed: Boolean, canStartCandidateEvaluation: Boolean, onStartRequirement: () -> Unit, onStopRequirement: () -> Unit, @@ -157,6 +158,7 @@ fun SearchProbeScreen( requirement = requirement, failureCode = requirementFailureCode, canStart = !active && + candidateCollectionAllowed && candidateEvaluationState != CandidateEvaluationState.RUNNING, onStart = onStartRequirement, onStop = onStopRequirement @@ -277,6 +279,7 @@ fun SearchProbeScreen( Button( onClick = onStart, enabled = readiness.canStartProbe && + candidateCollectionAllowed && requirementState != RequirementProbeState.RUNNING && candidateEvaluationState != CandidateEvaluationState.RUNNING, modifier = Modifier.fillMaxWidth(), @@ -338,6 +341,7 @@ private fun CandidateEvaluationSection( var rejectNote by remember(state) { mutableStateOf("") } val statusColor = when (state) { CandidateEvaluationState.AWAITING_CONFIRMATION, + CandidateEvaluationState.AWAITING_ADMIN_CONFIRMATION, CandidateEvaluationState.HUMAN_ACCEPTED -> colors.success CandidateEvaluationState.MANUAL_REVIEW, CandidateEvaluationState.NO_MATCH -> colors.warning @@ -575,6 +579,13 @@ private fun CandidateEvaluationSection( Text("拒绝本次候选") } } + CandidateEvaluationState.AWAITING_ADMIN_CONFIRMATION -> { + Text( + text = "候选已安全回传,等待后台采购员选择商品", + fontSize = 14.sp, + color = colors.success + ) + } CandidateEvaluationState.MANUAL_REVIEW, CandidateEvaluationState.NO_MATCH -> { if (candidateOrdinals.isNotEmpty()) { @@ -784,6 +795,7 @@ private fun candidateEvaluationStateLabel( CandidateEvaluationState.IDLE -> "等待" CandidateEvaluationState.RUNNING -> "评估中" CandidateEvaluationState.AWAITING_CONFIRMATION -> "等待人工确认" + CandidateEvaluationState.AWAITING_ADMIN_CONFIRMATION -> "等待后台确认" CandidateEvaluationState.MANUAL_REVIEW -> "需人工判断" CandidateEvaluationState.NO_MATCH -> "无建议候选" CandidateEvaluationState.HUMAN_ACCEPTED -> "人员已标记可用" diff --git a/android-buyer/app/src/main/java/com/roubao/autopilot/vlm/CandidateEvaluationModels.kt b/android-buyer/app/src/main/java/com/roubao/autopilot/vlm/CandidateEvaluationModels.kt index b2bc131..4c80ee5 100644 --- a/android-buyer/app/src/main/java/com/roubao/autopilot/vlm/CandidateEvaluationModels.kt +++ b/android-buyer/app/src/main/java/com/roubao/autopilot/vlm/CandidateEvaluationModels.kt @@ -142,6 +142,7 @@ enum class CandidateEvaluationState { IDLE, RUNNING, AWAITING_CONFIRMATION, + AWAITING_ADMIN_CONFIRMATION, MANUAL_REVIEW, NO_MATCH, HUMAN_ACCEPTED, diff --git a/android-buyer/app/src/test/java/com/roubao/autopilot/procurement/OrderCommandIntegrityTest.kt b/android-buyer/app/src/test/java/com/roubao/autopilot/procurement/OrderCommandIntegrityTest.kt new file mode 100644 index 0000000..ca0985b --- /dev/null +++ b/android-buyer/app/src/test/java/com/roubao/autopilot/procurement/OrderCommandIntegrityTest.kt @@ -0,0 +1,37 @@ +package com.roubao.autopilot.procurement + +import org.junit.Assert.assertEquals +import org.junit.Test + +class OrderCommandIntegrityTest { + @Test + fun hashMatchesGoContractVector() { + val command = PendingOrderCommand( + id = "00000000-0000-4000-8000-000000000009", + schemaVersion = 1, + type = "CREATE_PENDING_ORDER", + authorizationVersion = 3, + taskId = "00000000-0000-4000-8000-000000000001", + executionId = "00000000-0000-4000-8000-000000000002", + taskContentSha256 = "a".repeat(64), + originalSku = "红色 M", + quantity = 2, + candidate = OrderCommandCandidate( + candidateKey = "candidate-key-1", + observedOrdinal = 1, + title = "红色连衣裙", + skuText = "红色 / M", + priceText = "19.90 CNY", + cardSignature = "b".repeat(64), + detailSignature = "c".repeat(64), + detailEvidenceSha256 = "d".repeat(64), + specificationEvidenceSha256 = "e".repeat(64) + ), + commandSha256 = + "cf3cd802366d4de49d51940551385f7e5bd3283ef18cbf245a227c679d1058f3", + authorizationStatus = "DELIVERED" + ) + + assertEquals(command.commandSha256, OrderCommandIntegrity.sha256(command)) + } +} diff --git a/android-buyer/app/src/test/java/com/roubao/autopilot/procurement/OrderCommandSynchronizationTest.kt b/android-buyer/app/src/test/java/com/roubao/autopilot/procurement/OrderCommandSynchronizationTest.kt new file mode 100644 index 0000000..58dfaed --- /dev/null +++ b/android-buyer/app/src/test/java/com/roubao/autopilot/procurement/OrderCommandSynchronizationTest.kt @@ -0,0 +1,134 @@ +package com.roubao.autopilot.procurement + +import kotlinx.coroutines.runBlocking +import org.junit.Assert.assertEquals +import org.junit.Assert.assertTrue +import org.junit.Test + +class OrderCommandSynchronizationTest { + @Test + fun commandIsPersistedBeforeAcknowledgementAndReplayUsesSameKey() = runBlocking { + val task = task() + val execution = execution() + val delivered = command(task, execution) + val calls = mutableListOf() + var persisted: PendingOrderCommand? = null + + val result = OrderCommandSynchronization.synchronize( + existing = null, + task = task, + execution = execution, + pull = { + calls += "pull" + delivered + }, + persist = { + calls += if (it.acknowledged) "persist-ack" else "persist-command" + persisted = it + }, + acknowledge = { command, key -> + calls += "ack:$key" + assertEquals(command, persisted) + OrderCommandAcknowledgement( + commandId = command.id, + status = "ACKNOWLEDGED", + replayed = false + ) + }, + newIdempotencyKey = { "stable-ack-key" } + ) + + assertEquals( + listOf( + "pull", + "persist-command", + "ack:stable-ack-key", + "persist-ack" + ), + calls + ) + assertTrue(requireNotNull(result).acknowledged) + assertEquals("stable-ack-key", result.acknowledgementIdempotencyKey) + } + + @Test + fun persistenceFailurePreventsAcknowledgement() = runBlocking { + val task = task() + val execution = execution() + var acknowledged = false + + val failure = runCatching { + OrderCommandSynchronization.synchronize( + existing = null, + task = task, + execution = execution, + pull = { command(task, execution) }, + persist = { error("storage unavailable") }, + acknowledge = { _, _ -> + acknowledged = true + error("must not acknowledge") + }, + newIdempotencyKey = { "stable-ack-key" } + ) + }.exceptionOrNull() + + assertEquals("storage unavailable", failure?.message) + assertEquals(false, acknowledged) + } + + private fun task() = RemotePurchaseTask( + id = "00000000-0000-4000-8000-000000000001", + status = "WAITING_CONFIRMATION", + version = 4, + claimGeneration = 1, + claimExpiresAt = "2099-01-01T00:00:00Z", + title = "测试商品", + description = "", + sku = "红色 M", + imageAssetId = "00000000-0000-4000-8000-000000000003", + referenceImageUrl = "/api/v1/tasks/task/reference-image", + quantity = 2, + maxBudget = "20.00", + currency = "CNY" + ) + + private fun execution() = RunningExecution( + id = "00000000-0000-4000-8000-000000000002", + currentStep = "WAITING_ADMIN_CONFIRMATION", + expiresAt = "2099-01-01T00:00:00Z", + serverClockOffsetMillis = 0 + ) + + private fun command( + task: RemotePurchaseTask, + execution: RunningExecution + ): PendingOrderCommand { + val unsigned = PendingOrderCommand( + id = "00000000-0000-4000-8000-000000000009", + schemaVersion = 1, + type = "CREATE_PENDING_ORDER", + authorizationVersion = 1, + taskId = task.id, + executionId = execution.id, + taskContentSha256 = ExecutionTaskHash.sha256(task), + originalSku = task.sku, + quantity = task.quantity, + candidate = OrderCommandCandidate( + candidateKey = "f".repeat(64), + observedOrdinal = 1, + title = "测试候选", + skuText = "红色 / M", + priceText = "19.90", + cardSignature = "b".repeat(64), + detailSignature = "c".repeat(64), + detailEvidenceSha256 = "d".repeat(64), + specificationEvidenceSha256 = "e".repeat(64) + ), + commandSha256 = "0".repeat(64), + authorizationStatus = "DELIVERED" + ) + return unsigned.copy( + commandSha256 = OrderCommandIntegrity.sha256(unsigned) + ) + } +} diff --git a/android-buyer/app/src/test/java/com/roubao/autopilot/procurement/ProcurementApiClientTest.kt b/android-buyer/app/src/test/java/com/roubao/autopilot/procurement/ProcurementApiClientTest.kt index b052fff..46f88cc 100644 --- a/android-buyer/app/src/test/java/com/roubao/autopilot/procurement/ProcurementApiClientTest.kt +++ b/android-buyer/app/src/test/java/com/roubao/autopilot/procurement/ProcurementApiClientTest.kt @@ -250,6 +250,109 @@ class ProcurementApiClientTest { ) } + @Test + fun orderCommandPullAndAcknowledgementFollowDeviceContract() = runBlocking { + val execution = RunningExecution( + id = "execution-id", + currentStep = "WAITING_ADMIN_CONFIRMATION", + expiresAt = "2099-01-01T00:00:00Z", + serverClockOffsetMillis = 0 + ) + server.enqueue( + jsonResponse( + """ + { + "id":"command-id", + "schema_version":1, + "type":"CREATE_PENDING_ORDER", + "authorization_version":1, + "task_id":"task-id", + "execution_id":"execution-id", + "task_content_sha256":"${"a".repeat(64)}", + "original_sku":"SKU-01", + "quantity":2, + "candidate":{ + "candidate_key":"candidate-key-1", + "observed_ordinal":1, + "title":"候选商品", + "sku_text":"红色 / M", + "price_text":"20.00", + "card_signature":"${"b".repeat(64)}", + "detail_signature":"${"c".repeat(64)}", + "detail_evidence_sha256":"${"d".repeat(64)}", + "specification_evidence_sha256":"${"e".repeat(64)}" + }, + "command_sha256":"${"f".repeat(64)}", + "authorization_status":"DELIVERED" + } + """.trimIndent() + ) + ) + val command = requireNotNull( + api.pullOrderCommand( + session(), + task("/api/v1/tasks/task-id/reference-image?claim_generation=1"), + execution, + "claim-token-value" + ) + ) + val pullRequest = server.takeRequest() + assertEquals( + "/api/v1/tasks/task-id/commands/next", + pullRequest.path + ) + assertEquals( + "claim-token-value", + pullRequest.getHeader("X-Claim-Token") + ) + assertTrue(pullRequest.body.readUtf8().contains("\"execution_id\":\"execution-id\"")) + assertEquals("candidate-key-1", command.candidate.candidateKey) + + server.enqueue( + jsonResponse( + """{"command_id":"command-id","status":"ACKNOWLEDGED","replayed":false}""" + ) + ) + val acknowledgement = api.acknowledgeOrderCommand( + session(), + task("/api/v1/tasks/task-id/reference-image?claim_generation=1"), + execution, + "claim-token-value", + command, + "command-ack-idempotency-key" + ) + val ackRequest = server.takeRequest() + assertEquals( + "/api/v1/tasks/task-id/commands/command-id/ack", + ackRequest.path + ) + assertEquals( + "command-ack-idempotency-key", + ackRequest.getHeader("Idempotency-Key") + ) + assertTrue(ackRequest.body.readUtf8().contains(command.commandSha256)) + assertEquals("ACKNOWLEDGED", acknowledgement.status) + } + + @Test + fun noOrderCommandIsAStableEmptyResult() = runBlocking { + server.enqueue(MockResponse().setResponseCode(204)) + + val command = api.pullOrderCommand( + session(), + task("/api/v1/tasks/task-id/reference-image?claim_generation=1"), + RunningExecution( + id = "execution-id", + currentStep = "WAITING_ADMIN_CONFIRMATION", + expiresAt = "2099-01-01T00:00:00Z", + serverClockOffsetMillis = 0 + ), + "claim-token-value" + ) + + assertEquals(null, command) + } + private fun session() = ProcurementSession( backendUrl = server.url("/").toString().trimEnd('/'), username = "buyer01", diff --git a/backend-api/cmd/api/main.go b/backend-api/cmd/api/main.go index bfb8fc5..b15a23b 100644 --- a/backend-api/cmd/api/main.go +++ b/backend-api/cmd/api/main.go @@ -197,6 +197,14 @@ func buildRouter( if err != nil { return nil, err } + commands, err := usecase.NewDeviceOrderCommandService( + store, + clock, + ids, + ) + if err != nil { + return nil, err + } passwords, err := password.NewBcrypt(12) if err != nil { return nil, err @@ -206,6 +214,7 @@ func buildRouter( Lifecycle: lifecycle, Assets: assets, Results: results, + Commands: commands, }, ) if err != nil { diff --git a/backend-api/internal/domain/task.go b/backend-api/internal/domain/task.go index 4ec5c24..18afc6b 100644 --- a/backend-api/internal/domain/task.go +++ b/backend-api/internal/domain/task.go @@ -300,6 +300,32 @@ type OrderAuthorization struct { RevokedAt *time.Time FailureCode *string FailureMessage *string + CommandSHA256 *string + DeliveryAttemptCount int + LastDeliveredAt *time.Time +} + +type DeviceOrderCommand struct { + ID string + SchemaVersion int + Type string + AuthorizationVersion int + TaskID string + ExecutionID string + TaskContentSHA256 string + OriginalSKU string + Quantity int + CandidateKey string + ObservedOrdinal int + CandidateTitle string + CandidateSKUText string + CandidatePriceText string + CardSignature string + DetailSignature string + DetailEvidenceSHA256 string + SpecificationEvidenceSHA256 string + CommandSHA256 string + AuthorizationStatus OrderAuthorizationStatus } type ExecutionReport struct { diff --git a/backend-api/internal/platform/migration/claims_migration_test.go b/backend-api/internal/platform/migration/claims_migration_test.go index 1a2b453..3b6f09d 100644 --- a/backend-api/internal/platform/migration/claims_migration_test.go +++ b/backend-api/internal/platform/migration/claims_migration_test.go @@ -34,8 +34,11 @@ func TestClaimsMigrationPreservesHistoryAcrossUpDownUp(t *testing.T) { if applied, err := runner.Up(ctx); err != nil { t.Fatalf("initial Up() error = %v", err) - } else if applied != 8 { - t.Fatalf("initial Up() applied = %d, want 8", applied) + } else if applied != 9 { + t.Fatalf("initial Up() applied = %d, want 9", applied) + } + if err := runner.Down(ctx); err != nil { + t.Fatalf("initial Down(v9) error = %v", err) } if err := runner.Down(ctx); err != nil { t.Fatalf("initial Down(v8) error = %v", err) @@ -53,9 +56,14 @@ func TestClaimsMigrationPreservesHistoryAcrossUpDownUp(t *testing.T) { seedClaimsHistoricalFixture(t, db) if applied, err := runner.Up(ctx); err != nil { - t.Fatalf("Up(v5-v8) over historical data error = %v", err) - } else if applied != 4 { - t.Fatalf("Up(v5-v8) applied = %d, want 4", applied) + t.Fatalf("Up(v5-v9) over historical data error = %v", err) + } else if applied != 5 { + t.Fatalf("Up(v5-v9) applied = %d, want 5", applied) + } + assertClaimsHistory(t, db, true) + + if err := runner.Down(ctx); err != nil { + t.Fatalf("Down(v9) with compatible history error = %v", err) } assertClaimsHistory(t, db, true) @@ -85,9 +93,9 @@ func TestClaimsMigrationPreservesHistoryAcrossUpDownUp(t *testing.T) { assertClaimsHistory(t, db, false) if applied, err := runner.Up(ctx); err != nil { - t.Fatalf("final Up(v4-v8) error = %v", err) - } else if applied != 5 { - t.Fatalf("final Up(v4-v8) applied = %d, want 5", applied) + t.Fatalf("final Up(v4-v9) error = %v", err) + } else if applied != 6 { + t.Fatalf("final Up(v4-v9) applied = %d, want 6", applied) } assertClaimsHistory(t, db, true) } @@ -323,6 +331,9 @@ func TestClaimsMigrationDownFailsClosedForNewAuditData(t *testing.T) { t.Fatalf("insert v4 audit event: %v", err) } + if err := runner.Down(ctx); err != nil { + t.Fatalf("Down(v9) error = %v", err) + } if err := runner.Down(ctx); err != nil { t.Fatalf("Down(v8) error = %v", err) } diff --git a/backend-api/internal/platform/migration/runner_test.go b/backend-api/internal/platform/migration/runner_test.go index e4a0dc5..db65a2c 100644 --- a/backend-api/internal/platform/migration/runner_test.go +++ b/backend-api/internal/platform/migration/runner_test.go @@ -27,8 +27,8 @@ func TestRunnerSupportsUpStatusDownAndIdempotentUp(t *testing.T) { if err != nil { t.Fatalf("Up() error = %v", err) } - if applied != 8 { - t.Fatalf("Up() applied = %d, want 8", applied) + if applied != 9 { + t.Fatalf("Up() applied = %d, want 9", applied) } assertStatuses(t, runner, map[int64]bool{ 1: true, @@ -39,6 +39,7 @@ func TestRunnerSupportsUpStatusDownAndIdempotentUp(t *testing.T) { 6: true, 7: true, 8: true, + 9: true, }) applied, err = runner.Up(context.Background()) @@ -60,7 +61,8 @@ func TestRunnerSupportsUpStatusDownAndIdempotentUp(t *testing.T) { 5: true, 6: true, 7: true, - 8: false, + 8: true, + 9: false, }) applied, err = runner.Up(context.Background()) @@ -79,6 +81,7 @@ func TestRunnerSupportsUpStatusDownAndIdempotentUp(t *testing.T) { 6: true, 7: true, 8: true, + 9: true, }) } diff --git a/backend-api/internal/repository/sqlite/auth_repository_test.go b/backend-api/internal/repository/sqlite/auth_repository_test.go index d37a423..a6e0bc3 100644 --- a/backend-api/internal/repository/sqlite/auth_repository_test.go +++ b/backend-api/internal/repository/sqlite/auth_repository_test.go @@ -383,6 +383,9 @@ func TestAuthMigrationCanRollbackWithoutRebuildingPurchaseTasks( if err != nil { t.Fatalf("migration.New() error = %v", err) } + if err := runner.Down(context.Background()); err != nil { + t.Fatalf("Down(v9) error = %v", err) + } if err := runner.Down(context.Background()); err != nil { t.Fatalf("Down(v8) error = %v", err) } @@ -414,9 +417,9 @@ func TestAuthMigrationCanRollbackWithoutRebuildingPurchaseTasks( t.Fatal("purchase_tasks was lost during auth migration rollback") } if applied, err := runner.Up(context.Background()); err != nil { - t.Fatalf("Up(v3-v8) error = %v", err) - } else if applied != 6 { - t.Fatalf("Up(v3-v8) applied = %d, want 6", applied) + t.Fatalf("Up(v3-v9) error = %v", err) + } else if applied != 7 { + t.Fatalf("Up(v3-v9) applied = %d, want 7", applied) } } diff --git a/backend-api/internal/repository/sqlite/device_order_command_repository.go b/backend-api/internal/repository/sqlite/device_order_command_repository.go new file mode 100644 index 0000000..8b49156 --- /dev/null +++ b/backend-api/internal/repository/sqlite/device_order_command_repository.go @@ -0,0 +1,404 @@ +package sqlite + +import ( + "context" + "database/sql" + "errors" + "time" + + "cmroubao/backend-api/internal/domain" + "cmroubao/backend-api/internal/usecase" +) + +func (s *Store) PullDeviceOrderCommand( + ctx context.Context, + write usecase.PullDeviceOrderCommandWrite, +) (*domain.DeviceOrderCommand, error) { + tx, err := s.db.BeginTx(ctx, nil) + if err != nil { + return nil, repositoryFailure(err) + } + defer func() { _ = tx.Rollback() }() + + if err := validateDeviceOrderCommandClaim( + ctx, + tx, + write.UserID, + write.DeviceID, + write.TaskID, + write.ExecutionID, + write.ClaimGeneration, + write.ClaimTokenHash, + write.Now, + ); err != nil { + return nil, err + } + authorizationID, observedOrdinal, candidateTitle, err := + findActiveDeviceOrderAuthorization( + ctx, + tx, + write.TaskID, + write.ExecutionID, + ) + if errors.Is(err, sql.ErrNoRows) { + if err := tx.Commit(); err != nil { + return nil, repositoryFailure(err) + } + return nil, nil + } + if err != nil { + return nil, repositoryFailure(err) + } + authorization, err := getOrderAuthorization( + ctx, + tx, + authorizationID, + ) + if err != nil { + return nil, err + } + if authorization.UserID != write.UserID || + authorization.DeviceID != write.DeviceID || + authorization.ClaimGeneration != write.ClaimGeneration { + return nil, usecase.ErrExecutionMismatch + } + command := usecase.NewDeviceOrderCommand( + authorization, + observedOrdinal, + candidateTitle, + ) + commandHash, err := usecase.DeviceOrderCommandSHA256(command) + if err != nil { + return nil, repositoryFailure(err) + } + if authorization.CommandSHA256 != nil && + *authorization.CommandSHA256 != commandHash { + return nil, usecase.ErrRepositoryInvariant + } + switch authorization.Status { + case domain.OrderAuthorizationPendingDelivery: + result, err := tx.ExecContext( + ctx, + `UPDATE order_authorizations + SET status = 'DELIVERED', + command_sha256 = ?, + delivered_at = ?, + last_delivered_at = ?, + delivery_attempt_count = delivery_attempt_count + 1 + WHERE id = ? + AND status = 'PENDING_DELIVERY' + AND command_sha256 IS NULL`, + commandHash, + formatTimestamp(write.Now), + formatTimestamp(write.Now), + authorization.ID, + ) + if err != nil { + return nil, repositoryFailure(err) + } + affected, err := result.RowsAffected() + if err != nil { + return nil, repositoryFailure(err) + } + if affected != 1 { + return nil, usecase.ErrTaskStateConflict + } + if err := insertTaskEvent(ctx, tx, write.DeliveredEvent); err != nil { + return nil, err + } + command.AuthorizationStatus = domain.OrderAuthorizationDelivered + case domain.OrderAuthorizationDelivered, + domain.OrderAuthorizationAcknowledged: + result, err := tx.ExecContext( + ctx, + `UPDATE order_authorizations + SET last_delivered_at = ?, + delivery_attempt_count = delivery_attempt_count + 1 + WHERE id = ? + AND status IN ('DELIVERED', 'ACKNOWLEDGED') + AND command_sha256 = ?`, + formatTimestamp(write.Now), + authorization.ID, + commandHash, + ) + if err != nil { + return nil, repositoryFailure(err) + } + affected, err := result.RowsAffected() + if err != nil { + return nil, repositoryFailure(err) + } + if affected != 1 { + return nil, usecase.ErrTaskStateConflict + } + default: + return nil, usecase.ErrTaskStateConflict + } + command.CommandSHA256 = commandHash + if err := tx.Commit(); err != nil { + return nil, repositoryFailure(err) + } + return &command, nil +} + +func (s *Store) AcknowledgeDeviceOrderCommand( + ctx context.Context, + write usecase.AcknowledgeDeviceOrderCommandWrite, +) (domain.OrderAuthorization, bool, error) { + tx, err := s.db.BeginTx(ctx, nil) + if err != nil { + return domain.OrderAuthorization{}, false, repositoryFailure(err) + } + defer func() { _ = tx.Rollback() }() + + record, found, err := getDeviceOrderCommandAckRequest( + ctx, + tx, + write.DeviceID, + write.IdempotencyKey, + ) + if err != nil { + return domain.OrderAuthorization{}, false, err + } + if found { + if record.RequestSHA256 != write.RequestSHA256 || + record.TaskID != write.TaskID || + record.AuthorizationID != write.AuthorizationID { + return domain.OrderAuthorization{}, false, + usecase.ErrIdempotencyConflict + } + authorization, err := getOrderAuthorization( + ctx, + tx, + record.AuthorizationID, + ) + if err != nil { + return domain.OrderAuthorization{}, false, err + } + if err := tx.Commit(); err != nil { + return domain.OrderAuthorization{}, false, + repositoryFailure(err) + } + return authorization, true, nil + } + if err := validateDeviceOrderCommandClaim( + ctx, + tx, + write.UserID, + write.DeviceID, + write.TaskID, + write.ExecutionID, + write.ClaimGeneration, + write.ClaimTokenHash, + write.Now, + ); err != nil { + return domain.OrderAuthorization{}, false, err + } + authorization, err := getOrderAuthorization( + ctx, + tx, + write.AuthorizationID, + ) + if err != nil { + return domain.OrderAuthorization{}, false, err + } + if authorization.TaskID != write.TaskID || + authorization.ExecutionID != write.ExecutionID || + authorization.UserID != write.UserID || + authorization.DeviceID != write.DeviceID || + authorization.ClaimGeneration != write.ClaimGeneration { + return domain.OrderAuthorization{}, false, + usecase.ErrExecutionMismatch + } + if authorization.CommandSHA256 == nil || + *authorization.CommandSHA256 != write.CommandSHA256 { + return domain.OrderAuthorization{}, false, + usecase.ErrTaskStateConflict + } + switch authorization.Status { + case domain.OrderAuthorizationDelivered: + result, err := tx.ExecContext( + ctx, + `UPDATE order_authorizations + SET status = 'ACKNOWLEDGED', + acknowledged_at = ? + WHERE id = ? + AND status = 'DELIVERED' + AND command_sha256 = ?`, + formatTimestamp(write.Now), + authorization.ID, + write.CommandSHA256, + ) + if err != nil { + return domain.OrderAuthorization{}, false, + repositoryFailure(err) + } + affected, err := result.RowsAffected() + if err != nil { + return domain.OrderAuthorization{}, false, + repositoryFailure(err) + } + if affected != 1 { + return domain.OrderAuthorization{}, false, + usecase.ErrTaskStateConflict + } + if err := insertTaskEvent( + ctx, + tx, + write.AcknowledgedEvent, + ); err != nil { + return domain.OrderAuthorization{}, false, err + } + case domain.OrderAuthorizationAcknowledged: + default: + return domain.OrderAuthorization{}, false, + usecase.ErrTaskStateConflict + } + _, err = tx.ExecContext( + ctx, + `INSERT INTO device_order_command_ack_requests ( + device_id, idempotency_key, request_sha256, + task_id, authorization_id, created_at + ) VALUES (?, ?, ?, ?, ?, ?)`, + write.DeviceID, + write.IdempotencyKey, + write.RequestSHA256, + write.TaskID, + write.AuthorizationID, + formatTimestamp(write.Now), + ) + if err != nil { + return domain.OrderAuthorization{}, false, + repositoryFailure(err) + } + authorization, err = getOrderAuthorization( + ctx, + tx, + authorization.ID, + ) + if err != nil { + return domain.OrderAuthorization{}, false, err + } + if err := tx.Commit(); err != nil { + return domain.OrderAuthorization{}, false, repositoryFailure(err) + } + return authorization, false, nil +} + +func validateDeviceOrderCommandClaim( + ctx context.Context, + tx *sql.Tx, + userID string, + deviceID string, + taskID string, + executionID string, + claimGeneration int64, + claimTokenHash string, + now time.Time, +) error { + task, err := getClaimProtectedTask(ctx, tx, taskID) + if err != nil { + return err + } + if err := validateClaim( + task, + userID, + deviceID, + claimGeneration, + claimTokenHash, + now, + ); err != nil { + return err + } + if task.Status != domain.TaskStatusWaitingConfirmation || + task.CancelRequestedAt != nil { + return usecase.ErrTaskStateConflict + } + execution, err := getExecutionByID(ctx, tx, executionID) + if err != nil { + return err + } + if execution.TaskID != taskID || + execution.UserID != userID || + execution.DeviceID != deviceID || + execution.ClaimGeneration != claimGeneration || + execution.FinishedAt != nil { + return usecase.ErrExecutionMismatch + } + return nil +} + +type deviceOrderCommandAckRequest struct { + RequestSHA256 string + TaskID string + AuthorizationID string +} + +func getDeviceOrderCommandAckRequest( + ctx context.Context, + tx *sql.Tx, + deviceID string, + idempotencyKey string, +) (deviceOrderCommandAckRequest, bool, error) { + var record deviceOrderCommandAckRequest + err := tx.QueryRowContext( + ctx, + `SELECT request_sha256, task_id, authorization_id + FROM device_order_command_ack_requests + WHERE device_id = ? AND idempotency_key = ?`, + deviceID, + idempotencyKey, + ).Scan( + &record.RequestSHA256, + &record.TaskID, + &record.AuthorizationID, + ) + if errors.Is(err, sql.ErrNoRows) { + return deviceOrderCommandAckRequest{}, false, nil + } + if err != nil { + return deviceOrderCommandAckRequest{}, false, + repositoryFailure(err) + } + return record, true, nil +} + +func findActiveDeviceOrderAuthorization( + ctx context.Context, + tx *sql.Tx, + taskID string, + executionID string, +) (string, int, string, error) { + var authorizationID, candidateTitle string + var observedOrdinal int + err := tx.QueryRowContext( + ctx, + `SELECT + oa.id, + coi.candidate_ordinal, + co.title + FROM order_authorizations oa + JOIN candidate_observation_identities coi + ON coi.candidate_key = oa.candidate_key + AND coi.execution_id = oa.execution_id + JOIN candidate_observations co + ON co.execution_id = coi.execution_id + AND co.ordinal = coi.candidate_ordinal + WHERE oa.task_id = ? + AND oa.execution_id = ? + AND oa.status IN ( + 'PENDING_DELIVERY', + 'DELIVERED', + 'ACKNOWLEDGED' + ) + ORDER BY oa.authorization_version DESC + LIMIT 1`, + taskID, + executionID, + ).Scan( + &authorizationID, + &observedOrdinal, + &candidateTitle, + ) + return authorizationID, observedOrdinal, candidateTitle, err +} diff --git a/backend-api/internal/repository/sqlite/order_authorization_repository.go b/backend-api/internal/repository/sqlite/order_authorization_repository.go index a8c9f67..d994810 100644 --- a/backend-api/internal/repository/sqlite/order_authorization_repository.go +++ b/backend-api/internal/repository/sqlite/order_authorization_repository.go @@ -595,6 +595,7 @@ func scanOrderAuthorization( var supersedes, deliveredAt, acknowledgedAt sql.NullString var executionStartedAt, consumedAt, failedAt, revokedAt sql.NullString var failureCode, failureMessage sql.NullString + var commandSHA256, lastDeliveredAt sql.NullString var createdAt string err := scanner.Scan( &authorization.ID, @@ -629,6 +630,9 @@ func scanOrderAuthorization( &revokedAt, &failureCode, &failureMessage, + &commandSHA256, + &authorization.DeliveryAttemptCount, + &lastDeliveredAt, ) if err != nil { return domain.OrderAuthorization{}, err @@ -664,6 +668,13 @@ func scanOrderAuthorization( if failureMessage.Valid { authorization.FailureMessage = &failureMessage.String } + if commandSHA256.Valid { + authorization.CommandSHA256 = &commandSHA256.String + } + if authorization.LastDeliveredAt, err = + parseNullableTimestamp(lastDeliveredAt); err != nil { + return domain.OrderAuthorization{}, err + } return authorization, nil } @@ -676,7 +687,8 @@ const orderAuthorizationSelect = `SELECT specification_evidence_sha256, status, supersedes_authorization_id, created_by_user_id, created_at, delivered_at, acknowledged_at, execution_started_at, consumed_at, - failed_at, revoked_at, failure_code, failure_message + failed_at, revoked_at, failure_code, failure_message, + command_sha256, delivery_attempt_count, last_delivered_at FROM order_authorizations` const localAdminSubject = "local-admin" diff --git a/backend-api/internal/transport/httpapi/admin_handlers.go b/backend-api/internal/transport/httpapi/admin_handlers.go index 4762c0d..4065ba7 100644 --- a/backend-api/internal/transport/httpapi/admin_handlers.go +++ b/backend-api/internal/transport/httpapi/admin_handlers.go @@ -484,6 +484,9 @@ func orderAuthorizationResponse( "revoked_at": formatOptionalTime(authorization.RevokedAt), "failure_code": authorization.FailureCode, "failure_message": authorization.FailureMessage, + "command_sha256": authorization.CommandSHA256, + "delivery_attempt_count": authorization.DeliveryAttemptCount, + "last_delivered_at": formatOptionalTime(authorization.LastDeliveredAt), } } diff --git a/backend-api/internal/transport/httpapi/admin_handlers_test.go b/backend-api/internal/transport/httpapi/admin_handlers_test.go index 4766499..9c58a83 100644 --- a/backend-api/internal/transport/httpapi/admin_handlers_test.go +++ b/backend-api/internal/transport/httpapi/admin_handlers_test.go @@ -410,6 +410,9 @@ func TestAdminOrderAuthorizationIsIdempotentAndRevisioned(t *testing.T) { if err != nil { t.Fatalf("migration.New() error = %v", err) } + if err := runner.Down(context.Background()); err != nil { + t.Fatalf("device command migration down: %v", err) + } if err := runner.Down(context.Background()); err == nil { t.Fatal("order authorization migration down succeeded with retained data") } diff --git a/backend-api/internal/transport/httpapi/device_handlers.go b/backend-api/internal/transport/httpapi/device_handlers.go index a20875f..3a24b91 100644 --- a/backend-api/internal/transport/httpapi/device_handlers.go +++ b/backend-api/internal/transport/httpapi/device_handlers.go @@ -20,10 +20,12 @@ type DeviceServices struct { Lifecycle *usecase.LifecycleService Assets *usecase.AssetService Results *usecase.ExecutionResultService + Commands *usecase.DeviceOrderCommandService } func (services DeviceServices) validate() error { - if services.Lifecycle == nil || services.Assets == nil || services.Results == nil { + if services.Lifecycle == nil || services.Assets == nil || + services.Results == nil || services.Commands == nil { return errors.New("device services are required") } return nil @@ -73,12 +75,125 @@ func NewDeviceRouteRegistrar( routes.POST("/api/v1/tasks/:id/evidence", handler.uploadEvidence) routes.POST("/api/v1/tasks/:id/candidates", handler.storeCandidates) routes.POST("/api/v1/tasks/:id/human-reviews", handler.storeHumanReview) + routes.POST( + "/api/v1/tasks/:id/commands/next", + handler.pullOrderCommand, + ) + routes.POST( + "/api/v1/tasks/:id/commands/:command_id/ack", + handler.acknowledgeOrderCommand, + ) routes.POST("/api/v1/tasks/:id/complete", handler.completeTask) routes.POST("/api/v1/tasks/:id/fail", handler.failTask) return nil }, nil } +func (handler *deviceHandlers) pullOrderCommand(ctx *gin.Context) { + principal, ok := devicePrincipal(ctx) + if !ok { + return + } + var request struct { + DeviceID string `json:"device_id"` + ExecutionID string `json:"execution_id"` + ClaimGeneration int64 `json:"claim_generation"` + } + if !decodeDeviceJSON(ctx, &request) || + !deviceIDMatches(ctx, request.DeviceID, principal.DeviceID) { + return + } + command, err := handler.services.Commands.Pull( + ctx.Request.Context(), + usecase.PullDeviceOrderCommand{ + UserID: principal.UserID, + DeviceID: principal.DeviceID, + TaskID: ctx.Param("id"), + ExecutionID: request.ExecutionID, + ClaimGeneration: request.ClaimGeneration, + ClaimToken: ctx.GetHeader(claimTokenHeader), + }, + ) + if err != nil { + writeUsecaseError(ctx, err) + return + } + ctx.Header("Cache-Control", "no-store") + if command == nil { + ctx.Status(http.StatusNoContent) + return + } + ctx.JSON(http.StatusOK, deviceOrderCommandResponse(*command)) +} + +func (handler *deviceHandlers) acknowledgeOrderCommand(ctx *gin.Context) { + principal, ok := devicePrincipal(ctx) + if !ok { + return + } + var request struct { + DeviceID string `json:"device_id"` + ExecutionID string `json:"execution_id"` + ClaimGeneration int64 `json:"claim_generation"` + CommandSHA256 string `json:"command_sha256"` + } + if !decodeDeviceJSON(ctx, &request) || + !deviceIDMatches(ctx, request.DeviceID, principal.DeviceID) { + return + } + result, err := handler.services.Commands.Acknowledge( + ctx.Request.Context(), + usecase.AcknowledgeDeviceOrderCommand{ + UserID: principal.UserID, + DeviceID: principal.DeviceID, + TaskID: ctx.Param("id"), + ExecutionID: request.ExecutionID, + AuthorizationID: ctx.Param("command_id"), + ClaimGeneration: request.ClaimGeneration, + ClaimToken: ctx.GetHeader(claimTokenHeader), + CommandSHA256: request.CommandSHA256, + IdempotencyKey: ctx.GetHeader("Idempotency-Key"), + }, + ) + if err != nil { + writeUsecaseError(ctx, err) + return + } + ctx.Header("Cache-Control", "no-store") + ctx.JSON(http.StatusOK, gin.H{ + "command_id": result.Authorization.ID, + "status": result.Authorization.Status, + "replayed": result.Replayed, + }) +} + +func deviceOrderCommandResponse(command domain.DeviceOrderCommand) gin.H { + return gin.H{ + "id": command.ID, + "schema_version": command.SchemaVersion, + "type": command.Type, + "authorization_version": command.AuthorizationVersion, + "task_id": command.TaskID, + "execution_id": command.ExecutionID, + "task_content_sha256": command.TaskContentSHA256, + "original_sku": command.OriginalSKU, + "quantity": command.Quantity, + "candidate": gin.H{ + "candidate_key": command.CandidateKey, + "observed_ordinal": command.ObservedOrdinal, + "title": command.CandidateTitle, + "sku_text": command.CandidateSKUText, + "price_text": command.CandidatePriceText, + "card_signature": command.CardSignature, + "detail_signature": command.DetailSignature, + "detail_evidence_sha256": command.DetailEvidenceSHA256, + "specification_evidence_sha256": command.SpecificationEvidenceSHA256, + }, + "command_sha256": command.CommandSHA256, + "authorization_status": command.AuthorizationStatus, + } +} + func (handler *deviceHandlers) referenceImage(ctx *gin.Context) { principal, ok := devicePrincipal(ctx) if !ok { diff --git a/backend-api/internal/transport/httpapi/device_handlers_test.go b/backend-api/internal/transport/httpapi/device_handlers_test.go index 0737fde..057d421 100644 --- a/backend-api/internal/transport/httpapi/device_handlers_test.go +++ b/backend-api/internal/transport/httpapi/device_handlers_test.go @@ -713,9 +713,17 @@ func TestDeviceExecutionResultsAreIdempotentAndAuditable(t *testing.T) { if err != nil { t.Fatalf("migration.New() after review error = %v", err) } + if err := runner.Down(context.Background()); err != nil { + t.Fatalf("device command migration down: %v", err) + } if err := runner.Down(context.Background()); err == nil { t.Fatal("order workflow migration down succeeded with retained data") } + if applied, err := runner.Up(context.Background()); err != nil { + t.Fatalf("restore device command migration: %v", err) + } else if applied != 1 { + t.Fatalf("restored migrations = %d, want 1", applied) + } completePayload := fmt.Sprintf( `{"execution_id":%q,"claim_generation":%d,"task_content_sha256":%q,"execution_mode":"MANUAL_FIRST","outcome":"CANDIDATE_ACCEPTED","operator_reason":"人工核对标题、SKU和截图后接受","candidate":{"ordinal":1,"title":"手动候选","sku_text":"TEST-SKU","price":"12.00","product_url":"","image_url":"","card_signature":%q,"detail_signature":%q,"detail_evidence_sha256":%q,"specification_evidence_sha256":%q,"evidence_asset_ids":[%q,%q],"evaluation":null},"order_submitted":false}`, @@ -784,6 +792,326 @@ func TestDeviceExecutionResultsAreIdempotentAndAuditable(t *testing.T) { } } +func TestDeviceOrderCommandDeliveryAndAcknowledgementAreRecoverable( + t *testing.T, +) { + fixture := newDeviceHTTPFixture(t) + requireDeviceStatus(t, fixture.readyHeartbeat(t), http.StatusOK) + taskID := fixture.createPendingTask(t) + claim := fixture.claimNext(t, "order-command-claim", testOpaqueToken) + requireDeviceStatus(t, claim, http.StatusOK) + var claimed deviceLifecycleResponse + decodeResponse(t, claim, &claimed) + startPayload := fmt.Sprintf( + `{"claim_generation":%d,"expected_version":%d}`, + claimed.Task.ClaimGeneration, + claimed.Task.Version, + ) + start := performDeviceRequest(t, fixture.router, deviceRequest{ + method: http.MethodPost, + target: "/api/v1/tasks/" + taskID + "/start", + contentType: "application/json", + body: strings.NewReader(startPayload), + bearerToken: testOpaqueToken, + claimToken: testOpaqueToken, + idempotencyKey: "order-command-start", + }) + requireDeviceStatus(t, start, http.StatusOK) + var started deviceLifecycleResponse + decodeResponse(t, start, &started) + candidateKey := seedDeviceOrderCommandCandidate( + t, + fixture, + taskID, + started.Execution.ID, + ) + detail, err := fixture.tasks.Get( + context.Background(), + localAdminSubject, + taskID, + ) + if err != nil { + t.Fatalf("get waiting task detail: %v", err) + } + authorizations, err := usecase.NewOrderAuthorizationService( + fixture.store, + usecase.SystemClock{}, + usecase.UUIDGenerator{}, + ) + if err != nil { + t.Fatalf("NewOrderAuthorizationService() error = %v", err) + } + created, err := authorizations.Create( + context.Background(), + usecase.CreateOrderAuthorizationCommand{ + ActorUserID: deviceTestAdminID, + TaskID: taskID, + IdempotencyKey: "order-command-authorization", + ExecutionID: started.Execution.ID, + TaskContentSHA256: usecase.TaskContentSHA256(detail.Task), + ExpectedTaskVersion: detail.Task.Version, + CandidateKey: candidateKey, + ReasonSchemaVersion: 1, + PrimaryReasonCode: "SELECTED_BEST_MATCH", + Items: []usecase.OrderAuthorizationItemInput{{ + CandidateKey: candidateKey, + Label: "ACCEPT", + PrimaryReasonCode: "SKU_MATCH", + ReasonCodes: []string{"SKU_MATCH", "IMAGE_MATCH"}, + }}, + }, + ) + if err != nil { + t.Fatalf("create order authorization: %v", err) + } + pullPayload := fmt.Sprintf( + `{"device_id":%q,"execution_id":%q,"claim_generation":%d}`, + deviceTestDeviceID, + started.Execution.ID, + started.Task.ClaimGeneration, + ) + pull := performDeviceRequest(t, fixture.router, deviceRequest{ + method: http.MethodPost, + target: "/api/v1/tasks/" + taskID + "/commands/next", + contentType: "application/json", + body: strings.NewReader(pullPayload), + bearerToken: testOpaqueToken, + claimToken: testOpaqueToken, + }) + requireDeviceStatus(t, pull, http.StatusOK) + var command struct { + ID string `json:"id"` + Type string `json:"type"` + SchemaVersion int `json:"schema_version"` + TaskID string `json:"task_id"` + ExecutionID string `json:"execution_id"` + Quantity int `json:"quantity"` + CommandSHA256 string `json:"command_sha256"` + AuthorizationStatus string `json:"authorization_status"` + Candidate struct { + Key string `json:"candidate_key"` + ObservedOrdinal int `json:"observed_ordinal"` + Title string `json:"title"` + SKUText string `json:"sku_text"` + CardSignature string `json:"card_signature"` + DetailSignature string `json:"detail_signature"` + DetailEvidence string `json:"detail_evidence_sha256"` + SpecificationSHA string `json:"specification_evidence_sha256"` + } `json:"candidate"` + } + decodeResponse(t, pull, &command) + if command.ID != created.Authorization.ID || + command.Type != "CREATE_PENDING_ORDER" || + command.SchemaVersion != 1 || + command.TaskID != taskID || + command.ExecutionID != started.Execution.ID || + command.Quantity != 2 || + len(command.CommandSHA256) != 64 || + command.AuthorizationStatus != "DELIVERED" || + command.Candidate.Key != candidateKey || + command.Candidate.ObservedOrdinal != 1 || + command.Candidate.Title != "设备命令候选" || + command.Candidate.SKUText != "TEST-SKU-COMMAND" { + t.Fatalf("order command = %+v", command) + } + replayedPull := performDeviceRequest(t, fixture.router, deviceRequest{ + method: http.MethodPost, + target: "/api/v1/tasks/" + taskID + "/commands/next", + contentType: "application/json", + body: strings.NewReader(pullPayload), + bearerToken: testOpaqueToken, + claimToken: testOpaqueToken, + }) + requireDeviceStatus(t, replayedPull, http.StatusOK) + if !strings.Contains( + replayedPull.Body.String(), + `"command_sha256":"`+command.CommandSHA256+`"`, + ) { + t.Fatalf("replayed command = %s", replayedPull.Body.String()) + } + badAckPayload := fmt.Sprintf( + `{"device_id":%q,"execution_id":%q,"claim_generation":%d,"command_sha256":%q}`, + deviceTestDeviceID, + started.Execution.ID, + started.Task.ClaimGeneration, + strings.Repeat("0", 64), + ) + badAck := performDeviceRequest(t, fixture.router, deviceRequest{ + method: http.MethodPost, + target: "/api/v1/tasks/" + taskID + "/commands/" + command.ID + "/ack", + contentType: "application/json", + body: strings.NewReader(badAckPayload), + bearerToken: testOpaqueToken, + claimToken: testOpaqueToken, + idempotencyKey: "order-command-ack-bad", + }) + requireDeviceStatus(t, badAck, http.StatusConflict) + ackPayload := strings.Replace( + badAckPayload, + strings.Repeat("0", 64), + command.CommandSHA256, + 1, + ) + ack := performDeviceRequest(t, fixture.router, deviceRequest{ + method: http.MethodPost, + target: "/api/v1/tasks/" + taskID + "/commands/" + command.ID + "/ack", + contentType: "application/json", + body: strings.NewReader(ackPayload), + bearerToken: testOpaqueToken, + claimToken: testOpaqueToken, + idempotencyKey: "order-command-ack", + }) + requireDeviceStatus(t, ack, http.StatusOK) + if !strings.Contains(ack.Body.String(), `"status":"ACKNOWLEDGED"`) { + t.Fatalf("ack response = %s", ack.Body.String()) + } + ackReplay := performDeviceRequest(t, fixture.router, deviceRequest{ + method: http.MethodPost, + target: "/api/v1/tasks/" + taskID + "/commands/" + command.ID + "/ack", + contentType: "application/json", + body: strings.NewReader(ackPayload), + bearerToken: testOpaqueToken, + claimToken: testOpaqueToken, + idempotencyKey: "order-command-ack", + }) + requireDeviceStatus(t, ackReplay, http.StatusOK) + if !strings.Contains(ackReplay.Body.String(), `"replayed":true`) { + t.Fatalf("ack replay = %s", ackReplay.Body.String()) + } + acknowledgedPull := performDeviceRequest( + t, + fixture.router, + deviceRequest{ + method: http.MethodPost, + target: "/api/v1/tasks/" + taskID + "/commands/next", + contentType: "application/json", + body: strings.NewReader(pullPayload), + bearerToken: testOpaqueToken, + claimToken: testOpaqueToken, + }, + ) + requireDeviceStatus(t, acknowledgedPull, http.StatusOK) + if !strings.Contains( + acknowledgedPull.Body.String(), + `"authorization_status":"ACKNOWLEDGED"`, + ) { + t.Fatalf("acknowledged pull = %s", acknowledgedPull.Body.String()) + } + var deliveredEvents, acknowledgedEvents int + for eventType, target := range map[string]*int{ + "ORDER_AUTHORIZATION_DELIVERED": &deliveredEvents, + "ORDER_AUTHORIZATION_ACKNOWLEDGED": &acknowledgedEvents, + } { + if err := fixture.db.QueryRow( + `SELECT COUNT(*) FROM task_events + WHERE task_id = ? AND event_type = ?`, + taskID, + eventType, + ).Scan(target); err != nil { + t.Fatalf("count %s events: %v", eventType, err) + } + } + if deliveredEvents != 1 || acknowledgedEvents != 1 { + t.Fatalf( + "delivery/ack events = %d/%d", + deliveredEvents, + acknowledgedEvents, + ) + } + runner, err := migration.New(fixture.db) + if err != nil { + t.Fatalf("migration.New() error = %v", err) + } + if err := runner.Down(context.Background()); err == nil { + t.Fatal("device command migration down succeeded with command data") + } +} + +func seedDeviceOrderCommandCandidate( + t *testing.T, + fixture *deviceHTTPFixture, + taskID string, + executionID string, +) string { + t.Helper() + detail, err := fixture.tasks.Get( + context.Background(), + localAdminSubject, + taskID, + ) + if err != nil { + t.Fatalf("get running task detail: %v", err) + } + taskHash := usecase.TaskContentSHA256(detail.Task) + now := time.Now().UTC().Format(time.RFC3339Nano) + candidateKey := strings.Repeat("7", 64) + if _, err := fixture.db.Exec( + `INSERT INTO candidate_search_runs ( + execution_id, task_id, task_content_sha256, execution_mode, + search_query, started_at, received_at, observation_count, + collection_complete, received_after_execution_expiry + ) VALUES (?, ?, ?, 'MANUAL_FIRST', 'PDD_IMAGE_SEARCH', ?, ?, 1, 1, 0)`, + executionID, + taskID, + taskHash, + now, + now, + ); err != nil { + t.Fatalf("seed order command search run: %v", err) + } + if _, err := fixture.db.Exec( + `INSERT INTO candidate_observations ( + execution_id, task_id, ordinal, title, sku_text, price_text, + product_url, image_url, evidence_asset_ids_json, + collection_status, observed_at + ) VALUES (?, ?, 1, '设备命令候选', 'TEST-SKU-COMMAND', '21.50', + '', '', '[]', 'COMPLETE', ?)`, + executionID, + taskID, + now, + ); err != nil { + t.Fatalf("seed order command observation: %v", err) + } + if _, err := fixture.db.Exec( + `INSERT INTO candidate_observation_identities ( + candidate_key, execution_id, candidate_ordinal, + card_signature, detail_signature, detail_evidence_sha256, + specification_evidence_sha256, identity_version, created_at + ) VALUES (?, ?, 1, ?, ?, ?, ?, 1, ?)`, + candidateKey, + executionID, + strings.Repeat("8", 64), + strings.Repeat("9", 64), + strings.Repeat("a", 64), + strings.Repeat("b", 64), + now, + ); err != nil { + t.Fatalf("seed order command identity: %v", err) + } + if _, err := fixture.db.Exec( + `UPDATE purchase_tasks + SET status = 'WAITING_CONFIRMATION', + version = version + 1, + updated_at = ? + WHERE id = ? AND status = 'RUNNING'`, + now, + taskID, + ); err != nil { + t.Fatalf("move order command task to waiting: %v", err) + } + if _, err := fixture.db.Exec( + `UPDATE task_executions + SET current_step = 'WAITING_ADMIN_CONFIRMATION', + last_heartbeat_at = ? + WHERE id = ?`, + now, + executionID, + ); err != nil { + t.Fatalf("move order command execution to waiting: %v", err) + } + return candidateKey +} + func TestDeviceReleaseReturnsClaimedTaskToPending(t *testing.T) { fixture := newDeviceHTTPFixture(t) requireDeviceStatus(t, fixture.readyHeartbeat(t), http.StatusOK) @@ -1020,11 +1348,16 @@ func newDeviceHTTPFixture(t *testing.T) *deviceHTTPFixture { if err != nil { t.Fatalf("usecase.NewExecutionResultService() error = %v", err) } + commands, err := usecase.NewDeviceOrderCommandService(store, clock, ids) + if err != nil { + t.Fatalf("usecase.NewDeviceOrderCommandService() error = %v", err) + } deviceRoutes, err := NewDeviceRouteRegistrar( DeviceServices{ Lifecycle: lifecycle, Assets: assets, Results: results, + Commands: commands, }, ) if err != nil { diff --git a/backend-api/internal/usecase/device_order_command_service.go b/backend-api/internal/usecase/device_order_command_service.go new file mode 100644 index 0000000..c80686a --- /dev/null +++ b/backend-api/internal/usecase/device_order_command_service.go @@ -0,0 +1,336 @@ +package usecase + +import ( + "context" + "crypto/sha256" + "encoding/hex" + "encoding/json" + "errors" + "strings" + "time" + + "cmroubao/backend-api/internal/domain" +) + +const ( + deviceOrderCommandSchemaVersion = 1 + deviceOrderCommandType = "CREATE_PENDING_ORDER" +) + +type PullDeviceOrderCommand struct { + UserID string + DeviceID string + TaskID string + ExecutionID string + ClaimGeneration int64 + ClaimToken string +} + +type PullDeviceOrderCommandWrite struct { + PullDeviceOrderCommand + ClaimTokenHash string + Now time.Time + DeliveredEvent domain.TaskEvent +} + +type AcknowledgeDeviceOrderCommand struct { + UserID string + DeviceID string + TaskID string + ExecutionID string + AuthorizationID string + ClaimGeneration int64 + ClaimToken string + CommandSHA256 string + IdempotencyKey string +} + +type AcknowledgeDeviceOrderCommandWrite struct { + AcknowledgeDeviceOrderCommand + ClaimTokenHash string + RequestSHA256 string + Now time.Time + AcknowledgedEvent domain.TaskEvent +} + +type AcknowledgeDeviceOrderCommandResult struct { + Authorization domain.OrderAuthorization + Replayed bool +} + +type DeviceOrderCommandRepository interface { + PullDeviceOrderCommand( + context.Context, + PullDeviceOrderCommandWrite, + ) (*domain.DeviceOrderCommand, error) + AcknowledgeDeviceOrderCommand( + context.Context, + AcknowledgeDeviceOrderCommandWrite, + ) (domain.OrderAuthorization, bool, error) +} + +type DeviceOrderCommandService struct { + repository DeviceOrderCommandRepository + clock Clock + ids IDGenerator +} + +func NewDeviceOrderCommandService( + repository DeviceOrderCommandRepository, + clock Clock, + ids IDGenerator, +) (*DeviceOrderCommandService, error) { + if repository == nil || clock == nil || ids == nil { + return nil, errors.New("device order command service dependencies are required") + } + return &DeviceOrderCommandService{ + repository: repository, + clock: clock, + ids: ids, + }, nil +} + +func (service *DeviceOrderCommandService) Pull( + ctx context.Context, + command PullDeviceOrderCommand, +) (*domain.DeviceOrderCommand, error) { + command = normalizePullDeviceOrderCommand(command) + fields := lifecycleClaimFields( + command.UserID, + command.DeviceID, + command.TaskID, + command.ClaimGeneration, + command.ClaimToken, + ) + if !isUUID(command.ExecutionID) { + fields["execution_id"] = "must be a UUID" + } + if len(fields) > 0 { + return nil, invalidError( + "ORDER_COMMAND_PULL_INVALID", + "order command pull request is invalid", + fields, + ) + } + eventID, err := service.ids.NewID() + if err != nil { + return nil, internalLifecycleFailure(err) + } + now := service.clock.Now().UTC() + userID, deviceID := command.UserID, command.DeviceID + result, err := service.repository.PullDeviceOrderCommand( + ctx, + PullDeviceOrderCommandWrite{ + PullDeviceOrderCommand: command, + ClaimTokenHash: hashSecret(command.ClaimToken), + Now: now, + DeliveredEvent: domain.TaskEvent{ + ID: eventID, + TaskID: command.TaskID, + ActorUserID: &userID, + ActorDeviceID: &deviceID, + Type: "ORDER_AUTHORIZATION_DELIVERED", + Message: "order authorization delivered to device", + OccurredAt: now, + }, + }, + ) + if err != nil { + return nil, wrapLifecycleRepositoryError(err) + } + return result, nil +} + +func (service *DeviceOrderCommandService) Acknowledge( + ctx context.Context, + command AcknowledgeDeviceOrderCommand, +) (AcknowledgeDeviceOrderCommandResult, error) { + command = normalizeAcknowledgeDeviceOrderCommand(command) + fields := lifecycleClaimFields( + command.UserID, + command.DeviceID, + command.TaskID, + command.ClaimGeneration, + command.ClaimToken, + ) + if !isUUID(command.ExecutionID) { + fields["execution_id"] = "must be a UUID" + } + if !isUUID(command.AuthorizationID) { + fields["command_id"] = "must be a UUID" + } + if !sha256Pattern.MatchString(command.CommandSHA256) { + fields["command_sha256"] = "must be lowercase SHA-256" + } + validateIdempotencyField(fields, command.IdempotencyKey) + if len(fields) > 0 { + return AcknowledgeDeviceOrderCommandResult{}, invalidError( + "ORDER_COMMAND_ACK_INVALID", + "order command acknowledgement is invalid", + fields, + ) + } + requestSHA256, err := lifecycleRequestHash(struct { + UserID string `json:"user_id"` + DeviceID string `json:"device_id"` + TaskID string `json:"task_id"` + ExecutionID string `json:"execution_id"` + AuthorizationID string `json:"authorization_id"` + ClaimGeneration int64 `json:"claim_generation"` + ClaimTokenHash string `json:"claim_token_sha256"` + CommandSHA256 string `json:"command_sha256"` + }{ + UserID: command.UserID, + DeviceID: command.DeviceID, + TaskID: command.TaskID, + ExecutionID: command.ExecutionID, + AuthorizationID: command.AuthorizationID, + ClaimGeneration: command.ClaimGeneration, + ClaimTokenHash: hashSecret(command.ClaimToken), + CommandSHA256: command.CommandSHA256, + }) + if err != nil { + return AcknowledgeDeviceOrderCommandResult{}, + internalLifecycleFailure(err) + } + eventID, err := service.ids.NewID() + if err != nil { + return AcknowledgeDeviceOrderCommandResult{}, + internalLifecycleFailure(err) + } + now := service.clock.Now().UTC() + userID, deviceID := command.UserID, command.DeviceID + authorization, replayed, err := + service.repository.AcknowledgeDeviceOrderCommand( + ctx, + AcknowledgeDeviceOrderCommandWrite{ + AcknowledgeDeviceOrderCommand: command, + ClaimTokenHash: hashSecret(command.ClaimToken), + RequestSHA256: requestSHA256, + Now: now, + AcknowledgedEvent: domain.TaskEvent{ + ID: eventID, + TaskID: command.TaskID, + ActorUserID: &userID, + ActorDeviceID: &deviceID, + Type: "ORDER_AUTHORIZATION_ACKNOWLEDGED", + Message: "order authorization persisted by device", + OccurredAt: now, + }, + }, + ) + if err != nil { + return AcknowledgeDeviceOrderCommandResult{}, + wrapLifecycleRepositoryError(err) + } + return AcknowledgeDeviceOrderCommandResult{ + Authorization: authorization, + Replayed: replayed, + }, nil +} + +func DeviceOrderCommandSHA256( + command domain.DeviceOrderCommand, +) (string, error) { + payload := struct { + ID string `json:"id"` + SchemaVersion int `json:"schema_version"` + Type string `json:"type"` + AuthorizationVersion int `json:"authorization_version"` + TaskID string `json:"task_id"` + ExecutionID string `json:"execution_id"` + TaskContentSHA256 string `json:"task_content_sha256"` + OriginalSKU string `json:"original_sku"` + Quantity int `json:"quantity"` + Candidate struct { + CandidateKey string `json:"candidate_key"` + ObservedOrdinal int `json:"observed_ordinal"` + Title string `json:"title"` + SKUText string `json:"sku_text"` + PriceText string `json:"price_text"` + CardSignature string `json:"card_signature"` + DetailSignature string `json:"detail_signature"` + DetailEvidenceSHA256 string `json:"detail_evidence_sha256"` + SpecificationEvidenceSHA256 string `json:"specification_evidence_sha256"` + } `json:"candidate"` + }{ + ID: command.ID, + SchemaVersion: command.SchemaVersion, + Type: command.Type, + AuthorizationVersion: command.AuthorizationVersion, + TaskID: command.TaskID, + ExecutionID: command.ExecutionID, + TaskContentSHA256: command.TaskContentSHA256, + OriginalSKU: command.OriginalSKU, + Quantity: command.Quantity, + } + payload.Candidate.CandidateKey = command.CandidateKey + payload.Candidate.ObservedOrdinal = command.ObservedOrdinal + payload.Candidate.Title = command.CandidateTitle + payload.Candidate.SKUText = command.CandidateSKUText + payload.Candidate.PriceText = command.CandidatePriceText + payload.Candidate.CardSignature = command.CardSignature + payload.Candidate.DetailSignature = command.DetailSignature + payload.Candidate.DetailEvidenceSHA256 = command.DetailEvidenceSHA256 + payload.Candidate.SpecificationEvidenceSHA256 = + command.SpecificationEvidenceSHA256 + encoded, err := json.Marshal(payload) + if err != nil { + return "", err + } + sum := sha256.Sum256(encoded) + return hex.EncodeToString(sum[:]), nil +} + +func NewDeviceOrderCommand( + authorization domain.OrderAuthorization, + observedOrdinal int, + candidateTitle string, +) domain.DeviceOrderCommand { + return domain.DeviceOrderCommand{ + ID: authorization.ID, + SchemaVersion: deviceOrderCommandSchemaVersion, + Type: deviceOrderCommandType, + AuthorizationVersion: authorization.AuthorizationVersion, + TaskID: authorization.TaskID, + ExecutionID: authorization.ExecutionID, + TaskContentSHA256: authorization.TaskContentSHA256, + OriginalSKU: authorization.OriginalSKU, + Quantity: authorization.Quantity, + CandidateKey: authorization.CandidateKey, + ObservedOrdinal: observedOrdinal, + CandidateTitle: candidateTitle, + CandidateSKUText: authorization.CandidateSKUText, + CandidatePriceText: authorization.CandidatePriceText, + CardSignature: authorization.CardSignature, + DetailSignature: authorization.DetailSignature, + DetailEvidenceSHA256: authorization.DetailEvidenceSHA256, + SpecificationEvidenceSHA256: authorization.SpecificationEvidenceSHA256, + AuthorizationStatus: authorization.Status, + } +} + +func normalizePullDeviceOrderCommand( + command PullDeviceOrderCommand, +) PullDeviceOrderCommand { + command.UserID = strings.TrimSpace(command.UserID) + command.DeviceID = strings.TrimSpace(command.DeviceID) + command.TaskID = strings.TrimSpace(command.TaskID) + command.ExecutionID = strings.TrimSpace(command.ExecutionID) + command.ClaimToken = strings.TrimSpace(command.ClaimToken) + return command +} + +func normalizeAcknowledgeDeviceOrderCommand( + command AcknowledgeDeviceOrderCommand, +) AcknowledgeDeviceOrderCommand { + command.UserID = strings.TrimSpace(command.UserID) + command.DeviceID = strings.TrimSpace(command.DeviceID) + command.TaskID = strings.TrimSpace(command.TaskID) + command.ExecutionID = strings.TrimSpace(command.ExecutionID) + command.AuthorizationID = strings.TrimSpace(command.AuthorizationID) + command.ClaimToken = strings.TrimSpace(command.ClaimToken) + command.CommandSHA256 = strings.TrimSpace(command.CommandSHA256) + command.IdempotencyKey = strings.TrimSpace(command.IdempotencyKey) + return command +} diff --git a/backend-api/internal/usecase/device_order_command_service_test.go b/backend-api/internal/usecase/device_order_command_service_test.go new file mode 100644 index 0000000..6daaeb7 --- /dev/null +++ b/backend-api/internal/usecase/device_order_command_service_test.go @@ -0,0 +1,47 @@ +package usecase + +import ( + "testing" + + "cmroubao/backend-api/internal/domain" +) + +func TestDeviceOrderCommandSHA256MatchesMobileContractVector(t *testing.T) { + command := domain.DeviceOrderCommand{ + ID: "00000000-0000-4000-8000-000000000009", + SchemaVersion: 1, + Type: "CREATE_PENDING_ORDER", + AuthorizationVersion: 3, + TaskID: "00000000-0000-4000-8000-000000000001", + ExecutionID: "00000000-0000-4000-8000-000000000002", + TaskContentSHA256: repeatForCommandTest("a"), + OriginalSKU: "红色 M", + Quantity: 2, + CandidateKey: "candidate-key-1", + ObservedOrdinal: 1, + CandidateTitle: "红色连衣裙", + CandidateSKUText: "红色 / M", + CandidatePriceText: "19.90 CNY", + CardSignature: repeatForCommandTest("b"), + DetailSignature: repeatForCommandTest("c"), + DetailEvidenceSHA256: repeatForCommandTest("d"), + SpecificationEvidenceSHA256: repeatForCommandTest("e"), + } + + got, err := DeviceOrderCommandSHA256(command) + if err != nil { + t.Fatalf("DeviceOrderCommandSHA256() error = %v", err) + } + const want = "cf3cd802366d4de49d51940551385f7e5bd3283ef18cbf245a227c679d1058f3" + if got != want { + t.Fatalf("DeviceOrderCommandSHA256() = %q, want %q", got, want) + } +} + +func repeatForCommandTest(value string) string { + result := "" + for len(result) < 64 { + result += value + } + return result +} diff --git a/backend-api/migrations/00009_device_order_commands.sql b/backend-api/migrations/00009_device_order_commands.sql new file mode 100644 index 0000000..77b413f --- /dev/null +++ b/backend-api/migrations/00009_device_order_commands.sql @@ -0,0 +1,175 @@ +-- +goose Up +ALTER TABLE order_authorizations + ADD COLUMN command_sha256 TEXT + CHECK ( + command_sha256 IS NULL + OR ( + length(command_sha256) = 64 + AND command_sha256 NOT GLOB '*[^0-9a-f]*' + ) + ); + +ALTER TABLE order_authorizations + ADD COLUMN delivery_attempt_count INTEGER NOT NULL DEFAULT 0 + CHECK (delivery_attempt_count >= 0); + +ALTER TABLE order_authorizations + ADD COLUMN last_delivered_at TEXT; + +CREATE TABLE device_order_command_ack_requests ( + device_id TEXT NOT NULL + REFERENCES devices(id) ON UPDATE RESTRICT ON DELETE RESTRICT, + idempotency_key TEXT NOT NULL + CHECK ( + length(trim(idempotency_key)) > 0 + AND length(CAST(idempotency_key AS BLOB)) <= 128 + ), + request_sha256 TEXT NOT NULL + CHECK ( + length(request_sha256) = 64 + AND request_sha256 NOT GLOB '*[^0-9a-f]*' + ), + task_id TEXT NOT NULL + REFERENCES purchase_tasks(id) ON UPDATE RESTRICT ON DELETE RESTRICT, + authorization_id TEXT NOT NULL + REFERENCES order_authorizations(id) + ON UPDATE RESTRICT ON DELETE RESTRICT, + created_at TEXT NOT NULL, + PRIMARY KEY (device_id, idempotency_key) +); + +ALTER TABLE task_events RENAME TO task_events_v8; + +CREATE TABLE task_events ( + id TEXT PRIMARY KEY NOT NULL + CHECK (length(id) = 36), + task_id TEXT NOT NULL + REFERENCES purchase_tasks(id) ON UPDATE RESTRICT ON DELETE CASCADE, + event_type TEXT NOT NULL + CHECK ( + event_type IN ( + 'TASK_CREATED', + 'TASK_CLAIMED', + 'TASK_RECLAIMED', + 'TASK_RELEASED', + 'TASK_STARTED', + 'TASK_CANCEL_REQUESTED', + 'TASK_CANCELED', + 'CANDIDATES_READY', + 'ORDER_AUTHORIZATION_CREATED', + 'ORDER_AUTHORIZATION_DELIVERED', + 'ORDER_AUTHORIZATION_ACKNOWLEDGED' + ) + ), + message TEXT NOT NULL, + occurred_at TEXT NOT NULL, + actor_user_id TEXT + REFERENCES users(id) ON UPDATE RESTRICT ON DELETE RESTRICT, + actor_device_id TEXT + REFERENCES devices(id) ON UPDATE RESTRICT ON DELETE RESTRICT +); + +INSERT INTO task_events ( + id, task_id, event_type, message, occurred_at, + actor_user_id, actor_device_id +) +SELECT + id, task_id, event_type, message, occurred_at, + actor_user_id, actor_device_id +FROM task_events_v8; + +DROP TABLE task_events_v8; + +CREATE INDEX task_events_task_occurred_idx + ON task_events (task_id, occurred_at ASC, id ASC); + +CREATE INDEX task_events_actor_user_idx + ON task_events (actor_user_id, occurred_at DESC, id DESC); + +CREATE INDEX task_events_actor_device_idx + ON task_events (actor_device_id, occurred_at DESC, id DESC); + +-- +goose Down +CREATE TEMP TABLE device_order_commands_v9_down_guard ( + allowed INTEGER NOT NULL + CHECK (allowed = 1) +); + +INSERT INTO device_order_commands_v9_down_guard (allowed) +SELECT CASE + WHEN EXISTS (SELECT 1 FROM device_order_command_ack_requests) + OR EXISTS ( + SELECT 1 + FROM order_authorizations + WHERE command_sha256 IS NOT NULL + OR delivery_attempt_count > 0 + OR last_delivered_at IS NOT NULL + OR status IN ('DELIVERED', 'ACKNOWLEDGED') + ) + OR EXISTS ( + SELECT 1 + FROM task_events + WHERE event_type IN ( + 'ORDER_AUTHORIZATION_DELIVERED', + 'ORDER_AUTHORIZATION_ACKNOWLEDGED' + ) + ) + THEN 0 + ELSE 1 +END; + +DROP TABLE device_order_commands_v9_down_guard; +DROP TABLE device_order_command_ack_requests; + +ALTER TABLE task_events RENAME TO task_events_v9; + +CREATE TABLE task_events ( + id TEXT PRIMARY KEY NOT NULL + CHECK (length(id) = 36), + task_id TEXT NOT NULL + REFERENCES purchase_tasks(id) ON UPDATE RESTRICT ON DELETE CASCADE, + event_type TEXT NOT NULL + CHECK ( + event_type IN ( + 'TASK_CREATED', + 'TASK_CLAIMED', + 'TASK_RECLAIMED', + 'TASK_RELEASED', + 'TASK_STARTED', + 'TASK_CANCEL_REQUESTED', + 'TASK_CANCELED', + 'CANDIDATES_READY', + 'ORDER_AUTHORIZATION_CREATED' + ) + ), + message TEXT NOT NULL, + occurred_at TEXT NOT NULL, + actor_user_id TEXT + REFERENCES users(id) ON UPDATE RESTRICT ON DELETE RESTRICT, + actor_device_id TEXT + REFERENCES devices(id) ON UPDATE RESTRICT ON DELETE RESTRICT +); + +INSERT INTO task_events ( + id, task_id, event_type, message, occurred_at, + actor_user_id, actor_device_id +) +SELECT + id, task_id, event_type, message, occurred_at, + actor_user_id, actor_device_id +FROM task_events_v9; + +DROP TABLE task_events_v9; + +CREATE INDEX task_events_task_occurred_idx + ON task_events (task_id, occurred_at ASC, id ASC); + +CREATE INDEX task_events_actor_user_idx + ON task_events (actor_user_id, occurred_at DESC, id DESC); + +CREATE INDEX task_events_actor_device_idx + ON task_events (actor_device_id, occurred_at DESC, id DESC); + +ALTER TABLE order_authorizations DROP COLUMN last_delivered_at; +ALTER TABLE order_authorizations DROP COLUMN delivery_attempt_count; +ALTER TABLE order_authorizations DROP COLUMN command_sha256; diff --git a/docs/00-ai-start-here.md b/docs/00-ai-start-here.md index dddce04..43b27c2 100644 --- a/docs/00-ai-start-here.md +++ b/docs/00-ai-start-here.md @@ -58,7 +58,8 @@ T-205 原子领取/租约状态机、T-206 Android 登录/有限离线、T-207 证据回传、T-211 参考图召回和 SKU 硬匹配、T-212 候选身份映射,以及 T-213 受控 规格组合/价格核验均已完成。T-208 的原始候选观测、模型评估、确定性推荐、逐候选 结构化人工理由和修订历史也已完成。T-214 商品持久身份和重新定位指纹、T-215 -Admin 候选确认与不可变待投递授权也已完成;当前实现 T-216 设备命令投递与确认。 +Admin 候选确认、不可变待投递授权及 T-216 设备命令可靠投递与确认均已完成;下一项 +是 T-217 已授权商品重新定位与订单 dry-run。 不得直接把候选链接或列表 ordinal 当成授权。 手机从管理后端领取任务并回传结果,VLM、拼多多自动化和人工确认在 App 本地完成。 T-206 增加有限离线执行;T-207 已复用 Roubao 端上 OpenAI 兼容适配器并加密本地 Key。 diff --git a/docs/08-interaction-checklist.md b/docs/08-interaction-checklist.md index cf43e60..308d7c4 100644 --- a/docs/08-interaction-checklist.md +++ b/docs/08-interaction-checklist.md @@ -17,7 +17,7 @@ | IX-008 | US-006 | App/管理端错误状态 | 自动失败、取消、重试上传 | 显示结构化原因和恢复动作 | P0 | 已定 | | IX-009 | US-008 | App 候选理由/管理端决策详情 | 接受、拒绝、改选或修正 | 保存逐候选结构化人工标签 | P1 | T-208 已实现 | | IX-010 | US-009 | App 独立执行设置/同步状态 | 配置模式、离线执行或补报 | 授权内独立执行并可审计同步 | P0 | T-206 离线控制已实现,结果补报待 T-207 | -| IX-011 | US-010 | Admin 候选授权/App 订单执行 | 选择候选并授权、设备领取执行 | 创建一笔可对账的待付款订单并提醒人工付款 | P0 | T-215 已完成;T-216 进行中 | +| IX-011 | US-010 | Admin 候选授权/App 订单执行 | 选择候选并授权、设备领取执行 | 创建一笔可对账的待付款订单并提醒人工付款 | P0 | T-215/T-216 已完成;T-217 待实现 | ## IX-001 管理 Web 登录 diff --git a/docs/current-state.md b/docs/current-state.md index 53cee50..7c93ab7 100644 --- a/docs/current-state.md +++ b/docs/current-state.md @@ -5,7 +5,7 @@ ## 当前快照 - 日期:2026-07-28 -- 阶段:T-216 设备下单命令投递、确认与恢复进行中 +- 阶段:T-216 设备下单命令投递、确认与恢复已完成,待开始 T-217 - Git:当前分支为 `main`;T-001 至 T-004、T-101 至 T-104、T-201 至 T-207、T-209、 T-208、T-210、T-211、T-212、T-213、T-214 均已纳入 Git 历史 - 生产代码:`android-buyer/` 已接入 Roubao Android 源码 @@ -17,10 +17,10 @@ - 本机 Android 工具:JDK 17.0.13、Command-line Tools 22.0、SDK 34、 Build Tools 34.0.0、Platform Tools/ADB 37.0.0;用户级 SDK 环境变量已设置 - Android Studio:未安装;`winget` 静默安装卡住后已终止,不阻塞命令行构建 -- 测试:T-215 Android 单元测试、Debug/Release 构建和根 `init.ps1` 通过; - Debug APK `1.4.12 (17)` 已安装并启动于 PKG110 -- 后端测试:T-215 运行 `go test ./...`、`go test -race ./...`、`go vet ./...`; - v8 migration 往返、带授权数据的降级保护和 Admin 授权接口集成测试均通过 +- 测试:T-216 Android Debug/Release 单元测试与构建和根 `init.ps1` 通过; + Debug APK `1.4.13 (18)` 已构建但未覆盖安装,PKG110 当前仍为 `1.4.12 (17)` +- 后端测试:T-216 运行 `go test ./...`、`go test -race ./...`、`go vet ./...`; + v9 migration 往返、带命令数据的降级保护及 pull/ACK HTTP 集成测试均通过 - 原型:4 个管理 Web 页面和 7 个 Android 页面均可离线独立打开;Playwright 以 1440×900、390×844、360×800 验证 36 个页面/视口组合,无页面横向溢出、 脚本错误或外部请求,Android 可见交互控件均不小于 44px @@ -68,6 +68,10 @@ candidate key 选择并为全部候选填写结构化理由,服务端在同一事务追加人工 review、 不可变授权快照、事件和幂等记录。领取前可按版本改选,已投递后禁改;页面明确仅 授权创建设备端待付款订单,不授权付款。 +- T-216 设备命令:App 在 heartbeat 后主动拉取同一条不可变授权,严格校验跨端 + canonical hash、任务内容和候选四组指纹,先写 Keystore-backed 状态再幂等 ACK。 + 网络响应丢失和进程重启复用同一 ACK key;后台候选不再允许手机本地选择或重复 + 采集。本任务只同步授权,未打开拼多多或执行下单。 - Workflow:纯 Kotlin runner 已支持步骤 timeout、最多 3 次 retry、安全阻塞、 用户停止和单 runner 并发拒绝;T-102 已接入搜索加有界候选采集五步 - TaskSource:严格 CLI 已生成并验证真实私有 ProbeTask;默认 APK 不含私有 fixture @@ -81,7 +85,8 @@ - 离线执行:默认 30 分钟有限授权和 30 秒 best-effort heartbeat 已实现;PKG110 真机断开临时后端 95 秒后保持同一 execution,重连后滑动续期;到期持久安全停止, RUNNING 不自动重新分配 -- 测试设备:OnePlus PKG110,Android 16/API 36;肉包 `1.4.12 (17)`;拼多多 +- 测试设备:OnePlus PKG110,Android 16/API 36;已安装肉包 `1.4.12 (17)`, + 待安装构建为 `1.4.13 (18)`;拼多多 `8.17.0 (81700)` - 设备就绪:拼多多首页、文字/图片搜索、双列候选卡、详情截图和返回均已通过 8.17.0 真机验证;采购员已在 ColorOS 设置中手动启用肉包采购无障碍,APK 覆盖安装后授权 @@ -96,7 +101,7 @@ 已用 CLI 真实导入并逐字段/图片哈希验证,生成物位于被忽略的 `.local/` - 标准启动路径:`$env:RUN_START_COMMAND="1"; .\init.ps1` - 标准验证路径:`.\init.ps1` -- 当前 blocker:T-215 无阻塞。真实 VLM 服务地址、模型、设备级测试凭证、成本上限 +- 当前 blocker:T-216 无阻塞。真实 VLM 服务地址、模型、设备级测试凭证、成本上限 和数据留存尚未确认;当前只支持单 SKU/JPEG;候选探针截图要求 Android 11/API 30+ ## 当前目录 @@ -128,7 +133,7 @@ | `docs/tasks/T-208.md` | DONE | 归一化候选决策数据并增加结构化人工 review | | `docs/tasks/T-214.md` | DONE | 建立 execution-scoped candidate key 与设备采集指纹 | | `docs/tasks/T-215.md` | DONE | Admin 按 candidate key 选择并创建不可变待投递授权 | -| `docs/tasks/T-216.md` | DOING | App 主动拉取并先加密落盘再确认同一条下单命令 | +| `docs/tasks/T-216.md` | DONE | App 主动拉取并先加密落盘再确认同一条下单命令 | | `docs/design/` | 已确认 | T-202 原型索引、4 个管理页和 7 个 Android 页面 | | `deepseek总结.txt` | 已有 | 历史讨论摘要,不是正式需求权威 | | `android-buyer/` | 已有 | Roubao `main` 固定 commit 的 Android 基线 | @@ -139,10 +144,9 @@ ## 任务摘要 -- 已完成:T-001 至 T-004、T-101 至 T-104、T-201 至 T-215。 -- 正在进行:T-216 设备下单命令投递、确认与恢复。 -- 下一步:依次实现设备命令、订单 dry-run、单次提交对账和 - 付款提醒。 +- 已完成:T-001 至 T-004、T-101 至 T-104、T-201 至 T-216。 +- 正在进行:无。 +- 下一步:依次实现订单 dry-run、单次提交对账和付款提醒。 ## 当前可运行内容 diff --git a/docs/tasks/T-216.md b/docs/tasks/T-216.md index d66a137..27b0960 100644 --- a/docs/tasks/T-216.md +++ b/docs/tasks/T-216.md @@ -4,7 +4,7 @@ title: 设备下单命令投递、确认与恢复 phase: 2 deps: - T-215 -status: DOING +status: DONE created: 2026-07-28 context_ref: 827afc7 work_branch: null @@ -127,13 +127,13 @@ v9: ## 验收要点 -- [ ] 只有原 user/device/execution/generation/token 且租约有效时能拉取命令。 -- [ ] 首次拉取原子进入 `DELIVERED`;网络重放返回相同 command/hash 且不重复事件。 -- [ ] App 严格校验并先加密持久化,存储失败不 ACK,重启可恢复同一命令。 -- [ ] ACK 幂等进入 `ACKNOWLEDGED`,错误 hash/command/设备和并发请求被拒绝。 -- [ ] 后台候选流程等待 Admin,不再允许手机本地选择并提前完成任务。 -- [ ] App 等待/已授权状态可见,但不打开拼多多、不选择 SKU/数量、不提交订单。 -- [ ] v9 migration、Go test/race/vet、Android test/Debug/Release 和根验证通过。 +- [x] 只有原 user/device/execution/generation/token 且租约有效时能拉取命令。 +- [x] 首次拉取原子进入 `DELIVERED`;网络重放返回相同 command/hash 且不重复事件。 +- [x] App 严格校验并先加密持久化,存储失败不 ACK,重启可恢复同一命令。 +- [x] ACK 幂等进入 `ACKNOWLEDGED`,错误 hash/command/设备和并发请求被拒绝。 +- [x] 后台候选流程等待 Admin,不再允许手机本地选择并提前完成任务。 +- [x] App 等待/已授权状态可见,但不打开拼多多、不选择 SKU/数量、不提交订单。 +- [x] v9 migration、Go test/race/vet、Android test/Debug/Release 和根验证通过。 ## 边界 @@ -147,3 +147,13 @@ v9: - 2026-07-28:T-215 实现提交 `827afc7` 后领取。复用现有 30 秒前台 heartbeat 和 Keystore-backed 状态,选择 pull + durable-before-ACK,避免 Android 后台推送和 HTTP 不确定结果造成命令丢失。 +- 2026-07-28:新增 v9、命令 pull/ACK API、同一授权重放和投递/确认审计事件。 + HTTP 集成覆盖错误 hash、响应丢失后的相同 key 重试、ACK 后恢复拉取、事件不重复 + 及带数据禁止降级。 +- 2026-07-28:Android `1.4.13 (18)` 增加命令 canonical hash 校验、加密持久化、 + 固定 ACK key、重启恢复和等待/已授权 UI。后台任务有候选后不再出现手机接受/ + 拒绝入口,也禁止再次采集;空候选仍保留安全终态。 +- 2026-07-28:Go/Kotlin 固定哈希向量一致;Android 测试证明存储失败时 ACK 调用数 + 为 0。`go test ./...`、`go test -race ./...`、`go vet ./...`、Android `test`、 + Debug/Release assemble 和根 `.\init.ps1` 全部通过。APK 已构建,未在本任务中 + 覆盖安装或执行拼多多页面动作。